From 25b6cf6bcb13ae1602952d5655df098fab012125 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Krzysztof=20D=C4=85browski?= Date: Sat, 21 Mar 2026 16:19:20 +0100 Subject: [PATCH 01/12] add random subset button --- web/templates/training-datasets/new.html | 41 ++++++++++++++++++++++++ 1 file changed, 41 insertions(+) diff --git a/web/templates/training-datasets/new.html b/web/templates/training-datasets/new.html index 42d9fbf..a339538 100644 --- a/web/templates/training-datasets/new.html +++ b/web/templates/training-datasets/new.html @@ -92,6 +92,47 @@

Choose random subset

+
+
+
+ + +
+
+ + +
+
+ + +
+ +
+
+
From f6313de40cb87a3de95bc0f3ab8002cb456b4349 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Krzysztof=20D=C4=85browski?= Date: Mon, 13 Apr 2026 16:46:29 +0200 Subject: [PATCH 02/12] Update random subset flow and refactor training dataset UI --- web/templates/training-datasets/new.html | 95 ++++++++++++++++-------- 1 file changed, 65 insertions(+), 30 deletions(-) diff --git a/web/templates/training-datasets/new.html b/web/templates/training-datasets/new.html index a339538..9ce9806 100644 --- a/web/templates/training-datasets/new.html +++ b/web/templates/training-datasets/new.html @@ -94,43 +94,78 @@

Choose random subset

-
- - -
-
- - + hx-vals="js:{path: resolvePathForRequest(true)}"> +
+ +
-
- - +
+
+ + +
+
+ + +
+
+ + +
- + +
+ or + +
+
+ +
From 92cff6a47ccd3dcbc80e418e3fd5385073d7dde0 Mon Sep 17 00:00:00 2001 From: mytkom Date: Mon, 27 Apr 2026 12:21:50 +0200 Subject: [PATCH 03/12] add Sim to Data matching in AODBatcher --- cmd/aod_batcher/main.go | 348 +++++++++++++++--------- go.mod | 14 +- go.sum | 86 +++++- internal/config/config.go | 2 + internal/jalien/commands.go | 113 +++++--- internal/jalien/filepath_aod_matcher.go | 8 + internal/monalisa/scrape.go | 210 ++++++++++++++ internal/monalisa/scrape_test.go | 170 ++++++++++++ 8 files changed, 770 insertions(+), 181 deletions(-) create mode 100644 internal/monalisa/scrape.go create mode 100644 internal/monalisa/scrape_test.go diff --git a/cmd/aod_batcher/main.go b/cmd/aod_batcher/main.go index bb24817..3f3e54b 100644 --- a/cmd/aod_batcher/main.go +++ b/cmd/aod_batcher/main.go @@ -6,24 +6,30 @@ import ( "flag" "fmt" "log" - "math" "os" "path/filepath" "sort" "time" "github.com/mytkom/AliceTraINT/internal/jalien" + "github.com/mytkom/AliceTraINT/internal/monalisa" ) -type config struct { - path string - runs int - filesPerRun int - maxRunsPerBatch int - maxFilesPerBatch int - minSizeMB float64 - outputDir string +const MONALISA_DEFAULT_BASE_URL = "https://alimonitor.cern.ch" + +func getDataPath(dataTag, passName string, run uint64) string { + yearSuffix := dataTag[3:5] + return fmt.Sprintf("/alice/data/20%s/%s/%d/%s", yearSuffix, dataTag, run, passName) +} +type config struct { + path string + runs int + filesPerRun int + maxFilesPerBatch int + minSizeMB float64 + outputDir string + monalisaBaseUrl string jalienHost string jalienPort string clientCert string @@ -36,7 +42,6 @@ type metadata struct { Path string `json:"path"` Runs int `json:"runs"` FilesPerRun int `json:"files_per_run"` - MaxRunsPerBatch int `json:"max_runs_per_batch"` MaxFilesPerBatch int `json:"max_files_per_batch"` MinSizeMB float64 `json:"min_size_mb"` Timestamp string `json:"timestamp"` @@ -59,12 +64,12 @@ func parseFlags() (*config, error) { flag.StringVar(&cfg.path, "path", "", "JAliEn path under which to search for AOD files (e.g. /alice/sim/2024/LHC24f3)") flag.IntVar(&cfg.runs, "runs", 0, "Total number of runs to select") flag.IntVar(&cfg.filesPerRun, "files-per-run", 0, "Number of AOD files to select for each chosen run") - flag.IntVar(&cfg.maxRunsPerBatch, "max-runs-per-batch", -1, "Maximum number of runs to include in a single batch") flag.IntVar(&cfg.maxFilesPerBatch, "max-files-per-batch", -1, "Maximum number of AOD files to include in a single batch") flag.Float64Var(&cfg.minSizeMB, "min-size-mb", 0, "Optional minimal AOD file size in megabytes; files smaller than this are excluded") flag.StringVar(&cfg.outputDir, "output-dir", "", "Directory where batch .txt files will be written") flag.UintVar(&cfg.jalienTimeoutSeconds, "jalien-timeout-seconds", 600, "JAliEn timeout in seconds") // JAliEn connectivity; defaults from environment if flags are not provided. + flag.StringVar(&cfg.monalisaBaseUrl, "monalisa-base-url", MONALISA_DEFAULT_BASE_URL, "Monalisa base url used for querying anchor prod tag") flag.StringVar(&cfg.jalienHost, "jalien-host", os.Getenv("JALIEN_HOST"), "JAliEn host (default from $JALIEN_HOST or internal default)") flag.StringVar(&cfg.jalienPort, "jalien-port", os.Getenv("JALIEN_PORT"), "JAliEn port (default from $JALIEN_PORT or internal default)") flag.StringVar(&cfg.clientCert, "cert", os.Getenv("X509_USER_CERT"), "Path to user X.509 certificate (default from $X509_USER_CERT)") @@ -82,8 +87,8 @@ func parseFlags() (*config, error) { if cfg.filesPerRun <= 0 { return nil, errors.New("flag --files-per-run must be > 0") } - if cfg.maxRunsPerBatch < 0 && cfg.maxFilesPerBatch < 0 { - return nil, errors.New("flag --max-runs-per-batch or --max-files-per-batch must be > 0") + if cfg.maxFilesPerBatch <= 0 { + return nil, errors.New("flag --max-files-per-batch must be > 0") } if cfg.outputDir == "" { return nil, errors.New("flag --output-dir is required") @@ -100,32 +105,115 @@ func parseFlags() (*config, error) { if cfg.clientKey == "" { return nil, errors.New("flag --key (or $X509_USER_KEY) is required") } + if err := ensureOutputDirs(cfg.outputDir); err != nil { + return nil, err + } return cfg, nil } func run(cfg *config) error { - // Ensure output directory exists and is writable. - if err := os.MkdirAll(cfg.outputDir, 0o755); err != nil { - return fmt.Errorf("cannot create output directory %q: %w", cfg.outputDir, err) - } - client, err := jalien.NewClient(cfg.jalienHost, cfg.jalienPort, cfg.clientCert, cfg.clientKey, cfg.caCertsDir, cfg.jalienTimeoutSeconds) if err != nil { return fmt.Errorf("cannot create JAliEn client: %w", err) } - aods, err := client.FindAODFiles(cfg.path) + allSimAodsByRun, err := getSimAodsByRun(cfg.path, cfg.minSizeMB, client) if err != nil { - return fmt.Errorf("failed to discover AOD files under %q: %w", cfg.path, err) + return err + } + + simEligibleRuns, err := filterEligibleRuns(allSimAodsByRun, cfg.runs, cfg.filesPerRun) + if err != nil { + return err + } + + // Uniformly select runs + selectedRunIdx := uniformIndices(uint64(len(simEligibleRuns)), uint64(cfg.runs)) + selectedRuns := make([]uint64, 0, cfg.runs) + for _, idx := range selectedRunIdx { + selectedRuns = append(selectedRuns, simEligibleRuns[idx]) + } + + runToSimAods := make(map[uint64][]jalien.AODFile, len(selectedRuns)) + for _, rn := range selectedRuns { + // At this point we know len(files) >= filesPerRun by construction. + idxs := uniformIndices(uint64(len(allSimAodsByRun[rn])), uint64(cfg.filesPerRun)) + selected := make([]jalien.AODFile, 0, cfg.filesPerRun) + for _, idx := range idxs { + selected = append(selected, allSimAodsByRun[rn][idx]) + } + runToSimAods[rn] = selected + } + + // Get anchored data production (real data tag) + monc := monalisa.NewMonalisaClient(cfg.clientKey, cfg.clientCert, cfg.monalisaBaseUrl) + aodRow, err := monc.GetMCRow(allSimAodsByRun[selectedRuns[0]][0].LHCPeriod) + if err != nil { + return fmt.Errorf("error while getting anchor production tag: %s", err.Error()) + } + + runList, err := monc.GetRunList() + if err != nil { + return err + } + + runToDataAods := make(map[uint64][]jalien.AODFile, len(selectedRuns)) + for _, sr := range selectedRuns { + dataTag := "" + for _, t := range runList.RunToTags[sr] { + if !runList.TagToIsMC[t] { + dataTag = t + break + } + } + if dataTag == "" { + return fmt.Errorf("Cannot find a anchored production period for run %d", sr) + } + + dataPath := getDataPath(dataTag, aodRow.PassName, sr) + aods, err := getAods(dataPath, cfg.minSizeMB/2, client) + if err != nil { + return err + } + + fCount := min(cfg.filesPerRun, len(aods)) + idxs := uniformIndices(uint64(len(aods)), uint64(fCount)) + selected := make([]jalien.AODFile, 0, fCount) + for _, idx := range idxs { + selected = append(selected, aods[idx]) + } + + runToDataAods[sr] = selected + } + + if err := writeFilesInBatches(filepath.Join(cfg.outputDir, "sim"), flattenRunFiles(runToSimAods, selectedRuns), cfg.maxFilesPerBatch); err != nil { + return err + } + if err := writeFilesInBatches(filepath.Join(cfg.outputDir, "data"), flattenRunFiles(runToDataAods, selectedRuns), cfg.maxFilesPerBatch); err != nil { + return err + } + + // Persist non-sensitive configuration metadata alongside the batch files. + if err := writeMetadataFile(cfg); err != nil { + return err + } + + return nil +} + +func getAods(path string, minSizeMB float64, client *jalien.Client) ([]jalien.AODFile, error) { + aods, err := client.FindAODFiles(path) + if err != nil { + return nil, fmt.Errorf("failed to discover AOD files under %q: %w", path, err) } if len(aods) == 0 { - return fmt.Errorf("no AOD files found under path %q", cfg.path) + return nil, fmt.Errorf("no AOD files found under path %q", path) } // Apply optional minimum size filter (in MB) before any selection logic. - if cfg.minSizeMB > 0 { - minBytes := uint64(cfg.minSizeMB * 1024 * 1024) + if minSizeMB > 0 { + minBytes := uint64(minSizeMB * 1024 * 1024) filtered := make([]jalien.AODFile, 0, len(aods)) for _, f := range aods { if f.Size >= minBytes { @@ -133,17 +221,40 @@ func run(cfg *config) error { } } aods = filtered + + // sort for determinism + sort.Slice(aods, func(i, j int) bool { + if aods[i].AODNumber == aods[j].AODNumber { + return aods[i].Path < aods[j].Path + } + + return aods[i].AODNumber < aods[j].AODNumber + }) + if len(aods) == 0 { - return fmt.Errorf("no AOD files >= %.2f MB found under path %q", cfg.minSizeMB, cfg.path) + return nil, fmt.Errorf("no AOD files >= %.2f MB found under path %q", minSizeMB, path) } } + return aods, nil +} + +func getSimAodsByRun(path string, minSizeMB float64, client *jalien.Client) (map[uint64][]jalien.AODFile, error) { + aods, err := getAods(path, minSizeMB, client) + if err != nil { + return nil, err + } + // Group by run and sort. runsMap := make(map[uint64][]jalien.AODFile) for _, f := range aods { runsMap[f.RunNumber] = append(runsMap[f.RunNumber], f) } + return runsMap, nil +} + +func filterEligibleRuns(runsMap map[uint64][]jalien.AODFile, reqRunCount int, minFilesPerRun int) ([]uint64, error) { runNumbers := make([]uint64, 0, len(runsMap)) for rn := range runsMap { runNumbers = append(runNumbers, rn) @@ -154,129 +265,43 @@ func run(cfg *config) error { // not fail later when selecting files-per-run. eligibleRuns := make([]uint64, 0, len(runNumbers)) for _, rn := range runNumbers { - if len(runsMap[rn]) >= cfg.filesPerRun { + if len(runsMap[rn]) >= minFilesPerRun { eligibleRuns = append(eligibleRuns, rn) } } - if len(eligibleRuns) < cfg.runs { - return fmt.Errorf( - "requested %d runs with at least %d AOD files each, but only %d runs satisfy this under %q", - cfg.runs, cfg.filesPerRun, len(eligibleRuns), cfg.path, - ) - } - - selectedRunIdx := uniformIndices(len(eligibleRuns), cfg.runs) - selectedRuns := make([]uint64, 0, cfg.runs) - for _, idx := range selectedRunIdx { - selectedRuns = append(selectedRuns, eligibleRuns[idx]) - } - - // For determinism, within each selected run sort AODs by AODNumber then by path. - selectedFilesByRun := make(map[uint64][]jalien.AODFile, len(selectedRuns)) - for _, rn := range selectedRuns { - files := append([]jalien.AODFile(nil), runsMap[rn]...) - sort.Slice(files, func(i, j int) bool { - if files[i].AODNumber == files[j].AODNumber { - return files[i].Path < files[j].Path - } - return files[i].AODNumber < files[j].AODNumber - }) - - // At this point we know len(files) >= filesPerRun by construction. - idxs := uniformIndices(len(files), cfg.filesPerRun) - selected := make([]jalien.AODFile, 0, cfg.filesPerRun) - for _, idx := range idxs { - selected = append(selected, files[idx]) - } - selectedFilesByRun[rn] = selected - } - - // Build batches of runs. - if cfg.maxRunsPerBatch <= 0 && cfg.maxFilesPerBatch <= 0 { - return errors.New("flag --max-runs-per-batch or --max-files-per-batch must be > 0") - } - - if cfg.maxRunsPerBatch > 0 { - batchCount := int(math.Ceil(float64(len(selectedRuns)) / float64(cfg.maxRunsPerBatch))) - padding := len(fmt.Sprintf("%d", batchCount)) // number of digits for zero-padding - - batchIdx := 0 - for start := 0; start < len(selectedRuns); start += cfg.maxRunsPerBatch { - end := start + cfg.maxRunsPerBatch - if end > len(selectedRuns) { - end = len(selectedRuns) - } - batchRuns := selectedRuns[start:end] - batchFiles := make([]jalien.AODFile, 0, len(batchRuns)*cfg.filesPerRun) - for _, br := range batchRuns { - batchFiles = append(batchFiles, selectedFilesByRun[br]...) - } - - batchIdx++ - filename := fmt.Sprintf("batch_%0*d.txt", padding, batchIdx) - fullPath := filepath.Join(cfg.outputDir, filename) - - if err := writeBatchFile(fullPath, batchFiles); err != nil { - return err - } - } - } else if cfg.maxFilesPerBatch > 0 { - batchCount := int(math.Ceil(float64(len(selectedFilesByRun)) / float64(cfg.maxFilesPerBatch))) - padding := len(fmt.Sprintf("%d", batchCount)) // number of digits for zero-padding - files := make([]jalien.AODFile, 0, len(selectedFilesByRun)*cfg.filesPerRun) - for _, f := range selectedFilesByRun { - files = append(files, f...) - } - - batchIdx := 0 - for start := 0; start < len(files); start += cfg.maxFilesPerBatch { - end := start + cfg.maxFilesPerBatch - if end > len(files) { - end = len(files) - } - batchFiles := files[start:end] - - batchIdx++ - filename := fmt.Sprintf("batch_%0*d.txt", padding, batchIdx) - fullPath := filepath.Join(cfg.outputDir, filename) - - if err := writeBatchFile(fullPath, batchFiles); err != nil { - return err - } - } - } - // Persist non-sensitive configuration metadata alongside the batch files. - if err := writeMetadataFile(cfg); err != nil { - return err + if len(eligibleRuns) < reqRunCount { + return nil, fmt.Errorf( + "requested %d runs with at least %d AOD files each, but only %d runs satisfy this condition", + reqRunCount, minFilesPerRun, len(eligibleRuns)) } - return nil + return eligibleRuns, nil } // uniformIndices returns k indices in [0, n) that are as uniformly spaced as // possible over the interval when n >= k > 0. The result is strictly // increasing and deterministic. -func uniformIndices(n, k int) []int { +func uniformIndices(n, k uint64) []uint64 { if k <= 0 || n <= 0 { - return []int{} + return []uint64{} } if k == 1 { - return []int{n / 2} + return []uint64{n / 2} } step := float64(n) / float64(k) offset := step / 2.0 - indices := make([]int, 0, k) - prev := -1 - for i := 0; i < k; i++ { - x := int(offset + float64(i)*step) + indices := make([]uint64, 0, k) + prev := n + 1 + for i := uint64(0); i < k; i++ { + x := uint64(offset + float64(i)*step) if x >= n { x = n - 1 } // Ensure strict monotonicity in the unlikely case rounding produces duplicates. - if x <= prev { + if prev < n && x <= prev { x = prev + 1 if x >= n { x = n - 1 @@ -285,10 +310,36 @@ func uniformIndices(n, k int) []int { indices = append(indices, x) prev = x } + return indices } -func writeBatchFile(path string, batchFiles []jalien.AODFile) error { +func flattenRunFiles(runToAods map[uint64][]jalien.AODFile, runOrder []uint64) []jalien.AODFile { + flat := make([]jalien.AODFile, 0, len(runOrder)) + for _, rn := range runOrder { + flat = append(flat, runToAods[rn]...) + } + + return flat +} + +func writeFilesInBatches(outputDir string, files []jalien.AODFile, maxFilesPerBatch int) error { + batchCount := (len(files) + maxFilesPerBatch - 1) / maxFilesPerBatch + padding := len(fmt.Sprintf("%d", batchCount)) + + for start, batchIdx := 0, 1; start < len(files); start, batchIdx = start+maxFilesPerBatch, batchIdx+1 { + end := min(start+maxFilesPerBatch, len(files)) + filename := fmt.Sprintf("batch_%0*d.txt", padding, batchIdx) + if err := writeBatchFile(outputDir, filename, files[start:end]); err != nil { + return err + } + } + + return nil +} + +func writeBatchFile(output_dir, filename string, batchFiles []jalien.AODFile) error { + path := filepath.Join(output_dir, filename) f, err := os.Create(path) if err != nil { return fmt.Errorf("cannot create batch file %q: %w", path, err) @@ -308,12 +359,12 @@ func writeBatchFile(path string, batchFiles []jalien.AODFile) error { func writeMetadataFile(cfg *config) error { meta := metadata{ - Path: cfg.path, - Runs: cfg.runs, - FilesPerRun: cfg.filesPerRun, - MaxRunsPerBatch: cfg.maxRunsPerBatch, - MinSizeMB: cfg.minSizeMB, - Timestamp: time.Now().Format(time.RFC3339), + Path: cfg.path, + Runs: cfg.runs, + FilesPerRun: cfg.filesPerRun, + MaxFilesPerBatch: cfg.maxFilesPerBatch, + MinSizeMB: cfg.minSizeMB, + Timestamp: time.Now().Format(time.RFC3339), } path := filepath.Join(cfg.outputDir, "metadata.json") @@ -333,3 +384,38 @@ func writeMetadataFile(cfg *config) error { return nil } + +func ensureOutputDirs(outputDir string) error { + // 0o750: owner rwx, group rx, others no access. Keeps output private while + // still allowing read/execute traversal for shared project group. + for _, dir := range []string{ + outputDir, + filepath.Join(outputDir, "sim"), + filepath.Join(outputDir, "data"), + } { + if err := os.MkdirAll(dir, 0o750); err != nil { + return fmt.Errorf("cannot prepare output directory %q: %w", dir, err) + } + } + + metadataPath := filepath.Join(outputDir, "metadata.json") + testFile, err := os.CreateTemp(outputDir, ".write-check-*") + if err != nil { + return fmt.Errorf("output directory %q not writable: %w", outputDir, err) + } + testPath := testFile.Name() + _ = testFile.Close() + _ = os.Remove(testPath) + // Ensure we can replace/create metadata file later. + if _, err := os.Stat(metadataPath); err == nil { + f, err := os.OpenFile(metadataPath, os.O_WRONLY, 0) + if err != nil { + return fmt.Errorf("metadata file %q not writable: %w", metadataPath, err) + } + _ = f.Close() + } else if !errors.Is(err, os.ErrNotExist) { + return fmt.Errorf("cannot stat metadata file %q: %w", metadataPath, err) + } + + return nil +} diff --git a/go.mod b/go.mod index 5f58a44..86d9827 100644 --- a/go.mod +++ b/go.mod @@ -1,16 +1,17 @@ module github.com/mytkom/AliceTraINT -go 1.22.5 +go 1.25.0 require ( github.com/DATA-DOG/go-sqlmock v1.5.2 + github.com/PuerkitoBio/goquery v1.12.0 github.com/coreos/go-oidc v2.2.1+incompatible github.com/google/uuid v1.6.0 github.com/joho/godotenv v1.5.1 github.com/russross/blackfriday/v2 v2.1.0 github.com/stretchr/testify v1.9.0 github.com/thomasdarimont/go-kc-example v0.0.0-20170529223628-e3951d8faa4c - golang.org/x/crypto v0.26.0 + golang.org/x/crypto v0.49.0 golang.org/x/oauth2 v0.22.0 gorm.io/driver/postgres v1.5.9 gorm.io/driver/sqlite v1.5.6 @@ -19,8 +20,8 @@ require ( ) require ( + github.com/andybalholm/cascadia v1.3.3 // indirect github.com/davecgh/go-spew v1.1.2-0.20180830191138-d8f796af33cc // indirect - github.com/google/go-cmp v0.6.0 // indirect github.com/jackc/pgpassfile v1.0.0 // indirect github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761 // indirect github.com/jackc/pgx/v5 v5.6.0 // indirect @@ -33,9 +34,10 @@ require ( github.com/pquerna/cachecontrol v0.2.0 // indirect github.com/rogpeppe/go-internal v1.12.0 // indirect github.com/stretchr/objx v0.5.2 // indirect - golang.org/x/sync v0.8.0 // indirect - golang.org/x/sys v0.24.0 // indirect - golang.org/x/text v0.17.0 // indirect + golang.org/x/net v0.52.0 // indirect + golang.org/x/sync v0.20.0 // indirect + golang.org/x/sys v0.42.0 // indirect + golang.org/x/text v0.35.0 // indirect gopkg.in/square/go-jose.v2 v2.6.0 // indirect gopkg.in/yaml.v3 v3.0.1 // indirect ) diff --git a/go.sum b/go.sum index ae6e413..67160a6 100644 --- a/go.sum +++ b/go.sum @@ -1,5 +1,9 @@ github.com/DATA-DOG/go-sqlmock v1.5.2 h1:OcvFkGmslmlZibjAjaHm3L//6LiuBgolP7OputlJIzU= github.com/DATA-DOG/go-sqlmock v1.5.2/go.mod h1:88MAG/4G7SMwSE3CeA0ZKzrT5CiOU3OJ+JlNzwDqpNU= +github.com/PuerkitoBio/goquery v1.12.0 h1:pAcL4g3WRXekcB9AU/y1mbKez2dbY2AajVhtkO8RIBo= +github.com/PuerkitoBio/goquery v1.12.0/go.mod h1:802ej+gV2y7bbIhOIoPY5sT183ZW0YFofScC4q/hIpQ= +github.com/andybalholm/cascadia v1.3.3 h1:AG2YHrzJIm4BZ19iwJ/DAua6Btl3IwJX+VI4kktS1LM= +github.com/andybalholm/cascadia v1.3.3/go.mod h1:xNd9bqTn98Ln4DwST8/nG+H0yuB8Hmgu1YHNnWw0GeA= github.com/coreos/go-oidc v2.2.1+incompatible h1:mh48q/BqXqgjVHpy2ZY7WnWAbenxRjsz9N1i1YxjHAk= github.com/coreos/go-oidc v2.2.1+incompatible/go.mod h1:CgnwVTmzoESiwO9qyAFEMiHoZ1nMCKZlZ9V6mm3/LKc= github.com/creack/pty v1.1.9/go.mod h1:oKZEueFk5CKHvIhNR5MUki03XCEU+Q6VDXinZuGJ33E= @@ -52,16 +56,82 @@ github.com/stretchr/testify v1.9.0 h1:HtqpIVDClZ4nwg75+f6Lvsy/wHu+3BoSGCbBAcpTsT github.com/stretchr/testify v1.9.0/go.mod h1:r2ic/lqez/lEtzL7wO/rwa5dbSLXVDPFyf8C91i36aY= github.com/thomasdarimont/go-kc-example v0.0.0-20170529223628-e3951d8faa4c h1:lcylEBfIkt5tgRrDyh7N7a0U28d1bS7WurGpRl1fxJw= github.com/thomasdarimont/go-kc-example v0.0.0-20170529223628-e3951d8faa4c/go.mod h1:6c9GBv/SvzOyfPD5fEqQl9BCKS9IyO56JWB8BMp0CAc= -golang.org/x/crypto v0.26.0 h1:RrRspgV4mU+YwB4FYnuBoKsUapNIL5cohGAmSH3azsw= -golang.org/x/crypto v0.26.0/go.mod h1:GY7jblb9wI+FOo5y8/S2oY4zWP07AkOJ4+jxCqdqn54= +github.com/yuin/goldmark v1.4.13/go.mod h1:6yULJ656Px+3vBD8DxQVa3kxgyrAnzto9xy5taEt/CY= +golang.org/x/crypto v0.0.0-20190308221718-c2843e01d9a2/go.mod h1:djNgcEr1/C05ACkg1iLfiJU5Ep61QUkGW8qpdssI0+w= +golang.org/x/crypto v0.0.0-20210921155107-089bfa567519/go.mod h1:GvvjBRRGRdwPK5ydBHafDWAxML/pGHZbMvKqRZ5+Abc= +golang.org/x/crypto v0.13.0/go.mod h1:y6Z2r+Rw4iayiXXAIxJIDAJ1zMW4yaTpebo8fPOliYc= +golang.org/x/crypto v0.19.0/go.mod h1:Iy9bg/ha4yyC70EfRS8jz+B6ybOBKMaSxLj6P6oBDfU= +golang.org/x/crypto v0.23.0/go.mod h1:CKFgDieR+mRhux2Lsu27y0fO304Db0wZe70UKqHu0v8= +golang.org/x/crypto v0.31.0/go.mod h1:kDsLvtWBEx7MV9tJOj9bnXsPbxwJQ6csT/x4KIN4Ssk= +golang.org/x/crypto v0.49.0 h1:+Ng2ULVvLHnJ/ZFEq4KdcDd/cfjrrjjNSXNzxg0Y4U4= +golang.org/x/crypto v0.49.0/go.mod h1:ErX4dUh2UM+CFYiXZRTcMpEcN8b/1gxEuv3nODoYtCA= +golang.org/x/mod v0.6.0-dev.0.20220419223038-86c51ed26bb4/go.mod h1:jJ57K6gSWd91VN4djpZkiMVwK6gcyfeH4XE8wZrZaV4= +golang.org/x/mod v0.8.0/go.mod h1:iBbtSCu2XBx23ZKBPSOrRkjjQPZFPuis4dIYUhu/chs= +golang.org/x/mod v0.12.0/go.mod h1:iBbtSCu2XBx23ZKBPSOrRkjjQPZFPuis4dIYUhu/chs= +golang.org/x/mod v0.15.0/go.mod h1:hTbmBsO62+eylJbnUtE2MGJUyE7QWk4xUqPFrRgJ+7c= +golang.org/x/mod v0.17.0/go.mod h1:hTbmBsO62+eylJbnUtE2MGJUyE7QWk4xUqPFrRgJ+7c= +golang.org/x/net v0.0.0-20190620200207-3b0461eec859/go.mod h1:z5CRVTTTmAJ677TzLLGU+0bjPO0LkuOLi4/5GtJWs/s= +golang.org/x/net v0.0.0-20210226172049-e18ecbb05110/go.mod h1:m0MpNAwzfU5UDzcl9v0D8zg8gWTRqZa9RBIspLL5mdg= +golang.org/x/net v0.0.0-20220722155237-a158d28d115b/go.mod h1:XRhObCWvk6IyKnWLug+ECip1KBveYUHfp+8e9klMJ9c= +golang.org/x/net v0.6.0/go.mod h1:2Tu9+aMcznHK/AK1HMvgo6xiTLG5rD5rZLDS+rp2Bjs= +golang.org/x/net v0.10.0/go.mod h1:0qNGK6F8kojg2nk9dLZ2mShWaEBan6FAoqfSigmmuDg= +golang.org/x/net v0.15.0/go.mod h1:idbUs1IY1+zTqbi8yxTbhexhEEk5ur9LInksu6HrEpk= +golang.org/x/net v0.21.0/go.mod h1:bIjVDfnllIU7BJ2DNgfnXvpSvtn8VRwhlsaeUTyUS44= +golang.org/x/net v0.25.0/go.mod h1:JkAGAh7GEvH74S6FOH42FLoXpXbE/aqXSrIQjXgsiwM= +golang.org/x/net v0.33.0/go.mod h1:HXLR5J+9DxmrqMwG9qjGCxZ+zKXxBru04zlTvWlWuN4= +golang.org/x/net v0.52.0 h1:He/TN1l0e4mmR3QqHMT2Xab3Aj3L9qjbhRm78/6jrW0= +golang.org/x/net v0.52.0/go.mod h1:R1MAz7uMZxVMualyPXb+VaqGSa3LIaUqk0eEt3w36Sw= golang.org/x/oauth2 v0.22.0 h1:BzDx2FehcG7jJwgWLELCdmLuxk2i+x9UDpSiss2u0ZA= golang.org/x/oauth2 v0.22.0/go.mod h1:XYTD2NtWslqkgxebSiOHnXEap4TF09sJSc7H1sXbhtI= -golang.org/x/sync v0.8.0 h1:3NFvSEYkUoMifnESzZl15y791HH1qU2xm6eCJU5ZPXQ= -golang.org/x/sync v0.8.0/go.mod h1:Czt+wKu1gCyEFDUtn0jG5QVvpJ6rzVqr5aXyt9drQfk= -golang.org/x/sys v0.24.0 h1:Twjiwq9dn6R1fQcyiK+wQyHWfaz/BJB+YIpzU/Cv3Xg= -golang.org/x/sys v0.24.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA= -golang.org/x/text v0.17.0 h1:XtiM5bkSOt+ewxlOE/aE/AKEHibwj/6gvWMl9Rsh0Qc= -golang.org/x/text v0.17.0/go.mod h1:BuEKDfySbSR4drPmRPG/7iBdf8hvFMuRexcpahXilzY= +golang.org/x/sync v0.0.0-20190423024810-112230192c58/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= +golang.org/x/sync v0.0.0-20220722155255-886fb9371eb4/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= +golang.org/x/sync v0.1.0/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= +golang.org/x/sync v0.3.0/go.mod h1:FU7BRWz2tNW+3quACPkgCx/L+uEAv1htQ0V83Z9Rj+Y= +golang.org/x/sync v0.6.0/go.mod h1:Czt+wKu1gCyEFDUtn0jG5QVvpJ6rzVqr5aXyt9drQfk= +golang.org/x/sync v0.7.0/go.mod h1:Czt+wKu1gCyEFDUtn0jG5QVvpJ6rzVqr5aXyt9drQfk= +golang.org/x/sync v0.10.0/go.mod h1:Czt+wKu1gCyEFDUtn0jG5QVvpJ6rzVqr5aXyt9drQfk= +golang.org/x/sync v0.20.0 h1:e0PTpb7pjO8GAtTs2dQ6jYa5BWYlMuX047Dco/pItO4= +golang.org/x/sync v0.20.0/go.mod h1:9xrNwdLfx4jkKbNva9FpL6vEN7evnE43NNNJQ2LF3+0= +golang.org/x/sys v0.0.0-20190215142949-d0b11bdaac8a/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY= +golang.org/x/sys v0.0.0-20201119102817-f84b799fce68/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= +golang.org/x/sys v0.0.0-20210615035016-665e8c7367d1/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= +golang.org/x/sys v0.0.0-20220520151302-bc2c85ada10a/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= +golang.org/x/sys v0.0.0-20220722155257-8c9f86f7a55f/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= +golang.org/x/sys v0.5.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= +golang.org/x/sys v0.8.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= +golang.org/x/sys v0.12.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= +golang.org/x/sys v0.17.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA= +golang.org/x/sys v0.20.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA= +golang.org/x/sys v0.28.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA= +golang.org/x/sys v0.42.0 h1:omrd2nAlyT5ESRdCLYdm3+fMfNFE/+Rf4bDIQImRJeo= +golang.org/x/sys v0.42.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw= +golang.org/x/telemetry v0.0.0-20240228155512-f48c80bd79b2/go.mod h1:TeRTkGYfJXctD9OcfyVLyj2J3IxLnKwHJR8f4D8a3YE= +golang.org/x/term v0.0.0-20201126162022-7de9c90e9dd1/go.mod h1:bj7SfCRtBDWHUb9snDiAeCFNEtKQo2Wmx5Cou7ajbmo= +golang.org/x/term v0.0.0-20210927222741-03fcf44c2211/go.mod h1:jbD1KX2456YbFQfuXm/mYQcufACuNUgVhRMnK/tPxf8= +golang.org/x/term v0.5.0/go.mod h1:jMB1sMXY+tzblOD4FWmEbocvup2/aLOaQEp7JmGp78k= +golang.org/x/term v0.8.0/go.mod h1:xPskH00ivmX89bAKVGSKKtLOWNx2+17Eiy94tnKShWo= +golang.org/x/term v0.12.0/go.mod h1:owVbMEjm3cBLCHdkQu9b1opXd4ETQWc3BhuQGKgXgvU= +golang.org/x/term v0.17.0/go.mod h1:lLRBjIVuehSbZlaOtGMbcMncT+aqLLLmKrsjNrUguwk= +golang.org/x/term v0.20.0/go.mod h1:8UkIAJTvZgivsXaD6/pH6U9ecQzZ45awqEOzuCvwpFY= +golang.org/x/term v0.27.0/go.mod h1:iMsnZpn0cago0GOrHO2+Y7u7JPn5AylBrcoWkElMTSM= +golang.org/x/text v0.3.0/go.mod h1:NqM8EUOU14njkJ3fqMW+pc6Ldnwhi/IjpwHt7yyuwOQ= +golang.org/x/text v0.3.3/go.mod h1:5Zoc/QRtKVWzQhOtBMvqHzDpF6irO9z98xDceosuGiQ= +golang.org/x/text v0.3.7/go.mod h1:u+2+/6zg+i71rQMx5EYifcz6MCKuco9NR6JIITiCfzQ= +golang.org/x/text v0.7.0/go.mod h1:mrYo+phRRbMaCq/xk9113O4dZlRixOauAjOtrjsXDZ8= +golang.org/x/text v0.9.0/go.mod h1:e1OnstbJyHTd6l/uOt8jFFHp6TRDWZR/bV3emEE/zU8= +golang.org/x/text v0.13.0/go.mod h1:TvPlkZtksWOMsz7fbANvkp4WM8x/WCo/om8BMLbz+aE= +golang.org/x/text v0.14.0/go.mod h1:18ZOQIKpY8NJVqYksKHtTdi31H5itFRjB5/qKTNYzSU= +golang.org/x/text v0.15.0/go.mod h1:18ZOQIKpY8NJVqYksKHtTdi31H5itFRjB5/qKTNYzSU= +golang.org/x/text v0.21.0/go.mod h1:4IBbMaMmOPCJ8SecivzSH54+73PCFmPWxNTLm+vZkEQ= +golang.org/x/text v0.35.0 h1:JOVx6vVDFokkpaq1AEptVzLTpDe9KGpj5tR4/X+ybL8= +golang.org/x/text v0.35.0/go.mod h1:khi/HExzZJ2pGnjenulevKNX1W67CUy0AsXcNubPGCA= +golang.org/x/tools v0.0.0-20180917221912-90fa682c2a6e/go.mod h1:n7NCudcB/nEzxVGmLbDWY5pfWTLqBcC2KZ6jyYvM4mQ= +golang.org/x/tools v0.0.0-20191119224855-298f0cb1881e/go.mod h1:b+2E5dAYhXwXZwtnZ6UAqBI28+e2cm9otk0dWdXHAEo= +golang.org/x/tools v0.1.12/go.mod h1:hNGJHUnrk76NpqgfD5Aqm5Crs+Hm0VOH/i9J2+nxYbc= +golang.org/x/tools v0.6.0/go.mod h1:Xwgl3UAJ/d3gWutnCtw505GrjyAbvKui8lOU390QaIU= +golang.org/x/tools v0.13.0/go.mod h1:HvlwmtVNQAhOuCjW7xxvovg8wbNq7LwfXh/k7wXUl58= +golang.org/x/tools v0.21.1-0.20240508182429-e35e4ccd0d2d/go.mod h1:aiJjzUbINMkxbQROHiO6hDPo2LHcIPhhQsa9DLh0yGk= +golang.org/x/xerrors v0.0.0-20190717185122-a985d3407aa7/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0= gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c h1:Hei/4ADfdWqJk1ZMxUNpqntNwaWcugrBjAiHlqqRiVk= gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c/go.mod h1:JHkPIbrfpd72SG/EVd6muEfDQjcINNoR0C8j2r3qZ4Q= diff --git a/internal/config/config.go b/internal/config/config.go index ddd6ede..479ae02 100644 --- a/internal/config/config.go +++ b/internal/config/config.go @@ -17,6 +17,7 @@ type Config struct { JalienPort string JalienCertCADir string JalienTimeoutSeconds uint + MonalisaBaseURL string CCDBBaseURL string CCDBUploadSubdir string CertPath string @@ -61,6 +62,7 @@ func LoadConfig() *Config { JalienPort: getEnv("JALIEN_WSPORT", defaultJalienPort), JalienCertCADir: getEnv("JALIEN_CERT_CA_DIR", ""), JalienTimeoutSeconds: getEnvAsUint("JALIEN_TIMEOUT_SECONDS", 60), + MonalisaBaseURL: getEnv("MONALISA_URL", "https://alimonitor.cern.ch"), CCDBBaseURL: getEnv("CCDB_URL", "http://ccdb-test.cern.ch:8080"), CCDBUploadSubdir: getEnv("CCDB_UPLOAD_SUBDIR", "/Users/m/mmytkows"), CertPath: getEnv("GRID_CERT_PATH", ""), diff --git a/internal/jalien/commands.go b/internal/jalien/commands.go index 7c167de..7e7e302 100644 --- a/internal/jalien/commands.go +++ b/internal/jalien/commands.go @@ -122,42 +122,83 @@ func (client *Client) FindAODFiles(path string) ([]AODFile, error) { return nil, err } - matcher := newAODMatcher() - aods := make([]AODFile, 0, len(rawResults)) - - for _, r := range rawResults { - // Skip directories if type information is present. - if t := getStringField(r, "type"); t != "" { - if strings.ToLower(t) != "f" { - continue - } - } - - aodPath := getStringField(r, "lfn") - - pathVariables, err := matcher.MatchAO2DPath(aodPath) - if err != nil { - print("Skipping file: %s, error: %s", aodPath, err.Error()) - continue - } - if pathVariables == nil { - continue - } - - size, err := getUint64Field(r, "size") - if err != nil { - return nil, err - } - - aods = append(aods, AODFile{ - Name: aodFilename, - Path: aodPath, - Size: size, - LHCPeriod: pathVariables.LHCPeriod, - RunNumber: pathVariables.RunNumber, - AODNumber: pathVariables.AODNumber, - }) - } + aods := make([]AODFile, 0, len(rawResults)) + + if strings.Contains(path, "/alice/data") { + // process data + matcher := newDataAODMatcher() + + for i, r := range rawResults { + // Skip directories if type information is present. + if t := getStringField(r, "type"); t != "" { + if strings.ToLower(t) != "f" { + continue + } + } + + aodPath := getStringField(r, "lfn") + + pathVariables, err := matcher.MatchAO2DPath(aodPath) + if err != nil { + print("Skipping file: %s, error: %s", aodPath, err.Error()) + continue + } + if pathVariables == nil { + continue + } + + size, err := getUint64Field(r, "size") + if err != nil { + return nil, err + } + + aods = append(aods, AODFile{ + Name: aodFilename, + Path: aodPath, + Size: size, + LHCPeriod: pathVariables.LHCPeriod, + RunNumber: pathVariables.RunNumber, + AODNumber: uint64(i), + }) + } + } else { + // process MC + matcher := newAODMatcher() + + for _, r := range rawResults { + // Skip directories if type information is present. + if t := getStringField(r, "type"); t != "" { + if strings.ToLower(t) != "f" { + continue + } + } + + aodPath := getStringField(r, "lfn") + + pathVariables, err := matcher.MatchAO2DPath(aodPath) + if err != nil { + print("Skipping file: %s, error: %s", aodPath, err.Error()) + continue + } + if pathVariables == nil { + continue + } + + size, err := getUint64Field(r, "size") + if err != nil { + return nil, err + } + + aods = append(aods, AODFile{ + Name: aodFilename, + Path: aodPath, + Size: size, + LHCPeriod: pathVariables.LHCPeriod, + RunNumber: pathVariables.RunNumber, + AODNumber: pathVariables.AODNumber, + }) + } + } return aods, nil } diff --git a/internal/jalien/filepath_aod_matcher.go b/internal/jalien/filepath_aod_matcher.go index fcc209b..e90db8a 100644 --- a/internal/jalien/filepath_aod_matcher.go +++ b/internal/jalien/filepath_aod_matcher.go @@ -15,6 +15,14 @@ type filepathAODMatcherResult struct { AODNumber uint64 } +func newDataAODMatcher() *filepathAODMatcher { + regexpPattern := `/(?PLHC[a-zA-Z0-9]+)/(?P\d+)/.*/(?P\d+)/AO2D\.root` + + return &filepathAODMatcher{ + regexpCompiled: regexp.MustCompile(regexpPattern), + } +} + func newAODMatcher() *filepathAODMatcher { regexpPattern := `/(?PLHC[a-zA-Z0-9]+)/(?:\d+/)?(?P\d+)/(?:AOD/)?(?P\d+)/AO2D\.root` diff --git a/internal/monalisa/scrape.go b/internal/monalisa/scrape.go new file mode 100644 index 0000000..8e78603 --- /dev/null +++ b/internal/monalisa/scrape.go @@ -0,0 +1,210 @@ +package monalisa + +import ( + "bytes" + "crypto/tls" + "encoding/json" + "fmt" + "log" + "net/http" + "slices" + "strconv" + "strings" + "time" + + "github.com/PuerkitoBio/goquery" +) + +const MC_ANCH_PASS_NAME_COL_ID = 5 +const MC_TAG_TARGET_STR = "alicemctag" +const MC_ANCH_TARGET_STR = "alicemcaprod" + +type HyperloopRunList struct { + RunToTags map[uint64][]string + TagToRuns map[string][]uint64 + TagToIsMC map[string]bool +} + +type hyperloopRunListRawEntry struct { + EntryType string `json:"type"` + RunListStringified string `json:"runlist"` + PeriodTag string `json:"period"` +} + +func (rl *HyperloopRunList) UnmarshalJSON(raw []byte) error { + var rawEntries []hyperloopRunListRawEntry + err := json.Unmarshal(raw, &rawEntries) + if err != nil { + return err + } + + rl.TagToIsMC = make(map[string]bool, len(rawEntries)) + rl.TagToRuns = make(map[string][]uint64, len(rawEntries)) + rl.RunToTags = make(map[uint64][]string, len(rawEntries)) + + for _, e := range rawEntries { + for rs := range strings.SplitSeq(e.RunListStringified, ",") { + rn, err := strconv.ParseUint(rs, 10, 64) + if err != nil { + return err + } + + // check if run is present before appending + if !slices.Contains(rl.TagToRuns[e.PeriodTag], rn) { + rl.TagToRuns[e.PeriodTag] = append(rl.TagToRuns[e.PeriodTag], rn) + } + + // check if tag is present before appending + if !slices.Contains(rl.RunToTags[rn], e.PeriodTag) { + rl.RunToTags[rn] = append(rl.RunToTags[rn], e.PeriodTag) + } + + rl.TagToIsMC[e.PeriodTag] = (e.EntryType == "MC") + } + } + + return nil +} + +type MonalisaClient struct { + baseURL string + cert tls.Certificate + runListExpireAfterMinutes int64 + runList *HyperloopRunList + runListObtainedAt time.Time +} + +func NewMonalisaClient(keyPath, certPath, baseUrl string) *MonalisaClient { + cert, err := tls.LoadX509KeyPair(certPath, keyPath) + if err != nil { + log.Fatalf("cannot create MonalisaCert: %s", err.Error()) + } + return &MonalisaClient{ + baseURL: baseUrl, + cert: cert, + runListExpireAfterMinutes: 1440, // 24 Hours + } +} + +type MCRow struct { + MCTag string + AnchorProdTag string + PassName string +} + +func (c *MonalisaClient) GetMCRow(mcPeriod string) (*MCRow, error) { + url := fmt.Sprintf("%s/MC/?details=0&prodName=%s$", c.baseURL, mcPeriod) + + body := &bytes.Buffer{} + + request, err := http.NewRequest("GET", url, body) + if err != nil { + return nil, fmt.Errorf("failed to create request: %w", err) + } + + request.Header.Set("Content-Type", "application/json") + + client := &http.Client{ + Transport: &http.Transport{ + TLSClientConfig: &tls.Config{ + Certificates: []tls.Certificate{c.cert}, + InsecureSkipVerify: true, + }, + }, + } + resp, err := client.Do(request) + if err != nil { + return nil, fmt.Errorf("failed to upload file: %w", err) + } + + //nolint:errcheck + defer resp.Body.Close() + + if resp.StatusCode != http.StatusOK { + return nil, fmt.Errorf("upload failed with status code %d", resp.StatusCode) + } + + doc, err := goquery.NewDocumentFromReader(resp.Body) + if err != nil { + return nil, err + } + + mcTag := doc.Find(fmt.Sprintf(".link[target=%s]", MC_TAG_TARGET_STR)). + First(). + Text() + + anchorProd := doc.Find(fmt.Sprintf(".link[target=%s]", MC_ANCH_TARGET_STR)). + First(). + Text() + + passName := doc.Find("tbody tr.table_row td"). + Get(MC_ANCH_PASS_NAME_COL_ID). + FirstChild. + Data + + // trim string + passName = strings.Trim(passName, " \n\t") + + // debug + // fmt.Printf("MC Tag: %s, Anchor prod: %s, pass name: %s", mcTag, anchorProd, passName) + + return &MCRow{ + AnchorProdTag: anchorProd, + MCTag: mcTag, + PassName: passName, + }, nil +} + +func (c *MonalisaClient) GetRunList() (*HyperloopRunList, error) { + expireTimestamp := int64(c.runListObtainedAt.Unix()) + (c.runListExpireAfterMinutes * int64(time.Minute)) + if c.runList != nil && + time.Now().Unix() <= expireTimestamp { + return c.runList, nil + } + + url := fmt.Sprintf("%s/alihyperloop-data/runlist/list-runlist.jsp?lists=runlists", c.baseURL) + + body := &bytes.Buffer{} + req, err := http.NewRequest("GET", url, body) + if err != nil { + return nil, fmt.Errorf("problem with obtaining run list: %s", err.Error()) + } + + client := &http.Client{ + Transport: &http.Transport{ + TLSClientConfig: &tls.Config{ + Certificates: []tls.Certificate{c.cert}, + InsecureSkipVerify: true, + }, + }, + } + + resp, err := client.Do(req) + if err != nil { + return nil, fmt.Errorf("failed to get run list from hyperloop page: %w", err) + } + + //nolint:errcheck + defer resp.Body.Close() + + if resp.StatusCode != http.StatusOK { + return nil, fmt.Errorf("get run list failed with status code %d", resp.StatusCode) + } + + rawJson := bytes.Buffer{} + _, err = rawJson.ReadFrom(resp.Body) + if err != nil { + return nil, fmt.Errorf("reading body of response failed %s", err.Error()) + } + + var runList HyperloopRunList + err = json.Unmarshal(rawJson.Bytes(), &runList) + if err != nil { + return nil, fmt.Errorf("unmarshalling response from JSON failed %s", err.Error()) + } + + c.runList = &runList + c.runListObtainedAt = time.Now() + + return c.runList, nil +} diff --git a/internal/monalisa/scrape_test.go b/internal/monalisa/scrape_test.go new file mode 100644 index 0000000..ed2aaf8 --- /dev/null +++ b/internal/monalisa/scrape_test.go @@ -0,0 +1,170 @@ +package monalisa + +import ( + "crypto/tls" + "fmt" + "net/http" + "net/http/httptest" + "testing" + "time" +) + +func TestHyperloopRunListUnmarshalJSON(t *testing.T) { + raw := []byte(`[ + {"type":"MC","runlist":"101,102,101","period":"LHC24m"}, + {"type":"DATA","runlist":"102,103","period":"LHC24data"} + ]`) + + var got HyperloopRunList + if err := got.UnmarshalJSON(raw); err != nil { + t.Fatalf("UnmarshalJSON() error = %v, want nil", err) + } + + if len(got.TagToRuns["LHC24m"]) != 2 { + t.Fatalf("TagToRuns[LHC24m] len = %d, want 2", len(got.TagToRuns["LHC24m"])) + } + if len(got.RunToTags[102]) != 2 { + t.Fatalf("RunToTags[102] len = %d, want 2", len(got.RunToTags[102])) + } + if !got.TagToIsMC["LHC24m"] { + t.Fatalf("TagToIsMC[LHC24m] = false, want true") + } + if got.TagToIsMC["LHC24data"] { + t.Fatalf("TagToIsMC[LHC24data] = true, want false") + } +} + +func TestHyperloopRunListUnmarshalJSONBadRunNumber(t *testing.T) { + raw := []byte(`[{"type":"MC","runlist":"123,nope","period":"LHC24m"}]`) + + var got HyperloopRunList + if err := got.UnmarshalJSON(raw); err == nil { + t.Fatalf("UnmarshalJSON() error = nil, want non-nil") + } +} + +func TestGetRunListCachesResult(t *testing.T) { + requestCount := 0 + server := httptest.NewTLSServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + requestCount++ + if r.URL.Path != "/alihyperloop-data/runlist/list-runlist.jsp" { + t.Fatalf("path = %q, want runlist endpoint", r.URL.Path) + } + _, _ = w.Write([]byte(`[{"type":"MC","runlist":"2001,2002","period":"LHC24m"}]`)) + })) + defer server.Close() + + c := &MonalisaClient{ + baseURL: server.URL, + cert: tls.Certificate{}, + runListExpireAfterMinutes: 60, + } + + first, err := c.GetRunList() + if err != nil { + t.Fatalf("GetRunList() first error = %v, want nil", err) + } + second, err := c.GetRunList() + if err != nil { + t.Fatalf("GetRunList() second error = %v, want nil", err) + } + + if first != second { + t.Fatalf("cached pointer mismatch: first=%p second=%p", first, second) + } + if requestCount != 1 { + t.Fatalf("requestCount = %d, want 1", requestCount) + } +} + +func TestGetRunListRefreshesWhenExpired(t *testing.T) { + requestCount := 0 + server := httptest.NewTLSServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + requestCount++ + _, _ = w.Write([]byte(fmt.Sprintf(`[{"type":"MC","runlist":"%d","period":"LHC24m"}]`, 3000+requestCount))) + })) + defer server.Close() + + c := &MonalisaClient{ + baseURL: server.URL, + cert: tls.Certificate{}, + runListExpireAfterMinutes: -1, + runListObtainedAt: time.Now(), + runList: &HyperloopRunList{ + RunToTags: map[uint64][]string{1: {"old"}}, + TagToRuns: map[string][]uint64{"old": {1}}, + TagToIsMC: map[string]bool{"old": true}, + }, + } + + first, err := c.GetRunList() + if err != nil { + t.Fatalf("GetRunList() first error = %v, want nil", err) + } + second, err := c.GetRunList() + if err != nil { + t.Fatalf("GetRunList() second error = %v, want nil", err) + } + + if requestCount != 2 { + t.Fatalf("requestCount = %d, want 2", requestCount) + } + if first == second { + t.Fatalf("expected refreshed pointer, got same pointer %p", first) + } +} + +func TestGetMCRowParsesHTML(t *testing.T) { + const html = ` + + LHC24m + LHC24anchor + + + + + + +
c0c1c2c3c4 passAOD
+` + + server := httptest.NewTLSServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + _, _ = w.Write([]byte(html)) + })) + defer server.Close() + + c := &MonalisaClient{ + baseURL: server.URL, + cert: tls.Certificate{}, + } + + row, err := c.GetMCRow("LHC24m") + if err != nil { + t.Fatalf("GetMCRow() error = %v, want nil", err) + } + if row.MCTag != "LHC24m" { + t.Fatalf("MCTag = %q, want %q", row.MCTag, "LHC24m") + } + if row.AnchorProdTag != "LHC24anchor" { + t.Fatalf("AnchorProdTag = %q, want %q", row.AnchorProdTag, "LHC24anchor") + } + if row.PassName != "passAOD" { + t.Fatalf("PassName = %q, want %q", row.PassName, "passAOD") + } +} + +func TestGetMCRowNon200(t *testing.T) { + server := httptest.NewTLSServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.WriteHeader(http.StatusInternalServerError) + })) + defer server.Close() + + c := &MonalisaClient{ + baseURL: server.URL, + cert: tls.Certificate{}, + } + + if _, err := c.GetMCRow("LHC24m"); err == nil { + t.Fatalf("GetMCRow() error = nil, want non-nil") + } +} From d98b130a8c93d3aaed821d10ec37ba32b87ad478 Mon Sep 17 00:00:00 2001 From: mytkom Date: Mon, 27 Apr 2026 18:53:19 +0200 Subject: [PATCH 04/12] fallback to hyperloop if not exist in monalisa (O-O and Ne-Ne issue) --- cmd/aod_batcher/main.go | 29 ++++++++++++++++++++++------- internal/monalisa/scrape.go | 12 ++++++++++++ 2 files changed, 34 insertions(+), 7 deletions(-) diff --git a/cmd/aod_batcher/main.go b/cmd/aod_batcher/main.go index 3f3e54b..de0152b 100644 --- a/cmd/aod_batcher/main.go +++ b/cmd/aod_batcher/main.go @@ -8,6 +8,7 @@ import ( "log" "os" "path/filepath" + "regexp" "sort" "time" @@ -146,18 +147,32 @@ func run(cfg *config) error { runToSimAods[rn] = selected } - // Get anchored data production (real data tag) monc := monalisa.NewMonalisaClient(cfg.clientKey, cfg.clientCert, cfg.monalisaBaseUrl) - aodRow, err := monc.GetMCRow(allSimAodsByRun[selectedRuns[0]][0].LHCPeriod) - if err != nil { - return fmt.Errorf("error while getting anchor production tag: %s", err.Error()) - } - + // Get runlist runList, err := monc.GetRunList() if err != nil { return err } + // Get anchored data production (real data tag) + var passName string + mcTag := allSimAodsByRun[selectedRuns[0]][0].LHCPeriod + aodRow, err := monc.GetMCRow(mcTag) + if err != nil { + // try to obtain it from RunList + desc := runList.TagToDesc[mcTag] + re := regexp.MustCompile(`\bapass\d+(?:_\w+)?`) + m := re.FindString(desc) + + if m == "" { + return fmt.Errorf("error while obtaining anchored prod pass name") + } + + passName = m + } else { + passName = aodRow.PassName + } + runToDataAods := make(map[uint64][]jalien.AODFile, len(selectedRuns)) for _, sr := range selectedRuns { dataTag := "" @@ -171,7 +186,7 @@ func run(cfg *config) error { return fmt.Errorf("Cannot find a anchored production period for run %d", sr) } - dataPath := getDataPath(dataTag, aodRow.PassName, sr) + dataPath := getDataPath(dataTag, passName, sr) aods, err := getAods(dataPath, cfg.minSizeMB/2, client) if err != nil { return err diff --git a/internal/monalisa/scrape.go b/internal/monalisa/scrape.go index 8e78603..c15db64 100644 --- a/internal/monalisa/scrape.go +++ b/internal/monalisa/scrape.go @@ -23,12 +23,14 @@ type HyperloopRunList struct { RunToTags map[uint64][]string TagToRuns map[string][]uint64 TagToIsMC map[string]bool + TagToDesc map[string]string } type hyperloopRunListRawEntry struct { EntryType string `json:"type"` RunListStringified string `json:"runlist"` PeriodTag string `json:"period"` + Description string `json:"jt_description"` } func (rl *HyperloopRunList) UnmarshalJSON(raw []byte) error { @@ -39,6 +41,7 @@ func (rl *HyperloopRunList) UnmarshalJSON(raw []byte) error { } rl.TagToIsMC = make(map[string]bool, len(rawEntries)) + rl.TagToDesc = make(map[string]string, len(rawEntries)) rl.TagToRuns = make(map[string][]uint64, len(rawEntries)) rl.RunToTags = make(map[uint64][]string, len(rawEntries)) @@ -60,6 +63,7 @@ func (rl *HyperloopRunList) UnmarshalJSON(raw []byte) error { } rl.TagToIsMC[e.PeriodTag] = (e.EntryType == "MC") + rl.TagToDesc[e.PeriodTag] = e.Description } } @@ -133,10 +137,18 @@ func (c *MonalisaClient) GetMCRow(mcPeriod string) (*MCRow, error) { First(). Text() + if mcTag == "" { + return nil, fmt.Errorf("cannot obtain mcRow") + } + anchorProd := doc.Find(fmt.Sprintf(".link[target=%s]", MC_ANCH_TARGET_STR)). First(). Text() + if anchorProd == "" { + return nil, fmt.Errorf("cannot obtain mcRow") + } + passName := doc.Find("tbody tr.table_row td"). Get(MC_ANCH_PASS_NAME_COL_ID). FirstChild. From 9545b7ca99b85c71addc6542bb8dead6579f50e2 Mon Sep 17 00:00:00 2001 From: mytkom Date: Wed, 29 Apr 2026 17:19:18 +0200 Subject: [PATCH 05/12] divide data files per run number --- cmd/aod_batcher/main.go | 11 ++++++++++- 1 file changed, 10 insertions(+), 1 deletion(-) diff --git a/cmd/aod_batcher/main.go b/cmd/aod_batcher/main.go index de0152b..8df6f63 100644 --- a/cmd/aod_batcher/main.go +++ b/cmd/aod_batcher/main.go @@ -6,6 +6,7 @@ import ( "flag" "fmt" "log" + "math" "os" "path/filepath" "regexp" @@ -27,6 +28,7 @@ type config struct { path string runs int filesPerRun int + filesPerRunData int maxFilesPerBatch int minSizeMB float64 outputDir string @@ -43,6 +45,7 @@ type metadata struct { Path string `json:"path"` Runs int `json:"runs"` FilesPerRun int `json:"files_per_run"` + FilesPerRunData int `json:"files_per_run_data"` MaxFilesPerBatch int `json:"max_files_per_batch"` MinSizeMB float64 `json:"min_size_mb"` Timestamp string `json:"timestamp"` @@ -65,6 +68,7 @@ func parseFlags() (*config, error) { flag.StringVar(&cfg.path, "path", "", "JAliEn path under which to search for AOD files (e.g. /alice/sim/2024/LHC24f3)") flag.IntVar(&cfg.runs, "runs", 0, "Total number of runs to select") flag.IntVar(&cfg.filesPerRun, "files-per-run", 0, "Number of AOD files to select for each chosen run") + flag.IntVar(&cfg.filesPerRunData, "files-per-run-data", 0, "Number of AOD files to select for each chosen run's experimental data (typically 3-4x bigger files than for MC)") flag.IntVar(&cfg.maxFilesPerBatch, "max-files-per-batch", -1, "Maximum number of AOD files to include in a single batch") flag.Float64Var(&cfg.minSizeMB, "min-size-mb", 0, "Optional minimal AOD file size in megabytes; files smaller than this are excluded") flag.StringVar(&cfg.outputDir, "output-dir", "", "Directory where batch .txt files will be written") @@ -88,6 +92,10 @@ func parseFlags() (*config, error) { if cfg.filesPerRun <= 0 { return nil, errors.New("flag --files-per-run must be > 0") } + if cfg.filesPerRunData <= 0 { + cfg.filesPerRunData = int(math.Ceil(float64(cfg.filesPerRun) / 3.0)) + fmt.Printf("Setting up default files per run data as: %d (ceil(filesPerRun / 3))\n", cfg.filesPerRunData) + } if cfg.maxFilesPerBatch <= 0 { return nil, errors.New("flag --max-files-per-batch must be > 0") } @@ -192,7 +200,7 @@ func run(cfg *config) error { return err } - fCount := min(cfg.filesPerRun, len(aods)) + fCount := min(cfg.filesPerRunData, len(aods)) idxs := uniformIndices(uint64(len(aods)), uint64(fCount)) selected := make([]jalien.AODFile, 0, fCount) for _, idx := range idxs { @@ -377,6 +385,7 @@ func writeMetadataFile(cfg *config) error { Path: cfg.path, Runs: cfg.runs, FilesPerRun: cfg.filesPerRun, + FilesPerRunData: cfg.filesPerRunData, MaxFilesPerBatch: cfg.maxFilesPerBatch, MinSizeMB: cfg.minSizeMB, Timestamp: time.Now().Format(time.RFC3339), From df5c0ca5c1f9827d4aff0e43f216cf973830bf53 Mon Sep 17 00:00:00 2001 From: mytkom Date: Mon, 27 Jul 2026 12:59:43 +0200 Subject: [PATCH 06/12] update Dockerfile to new Golang version --- Dockerfile | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/Dockerfile b/Dockerfile index 3887cc4..d59c6ef 100644 --- a/Dockerfile +++ b/Dockerfile @@ -1,4 +1,4 @@ -FROM golang:1.22-alpine AS builder +FROM golang:1.25-alpine AS builder # --- config that differs locally vs OpenShift --- ARG CERT_PATH=./gridCertificate.p12 # default for local builds @@ -77,4 +77,4 @@ RUN mkdir -p . && \ USER 1001 EXPOSE 8088 -CMD ["./AliceTraINT"] \ No newline at end of file +CMD ["./AliceTraINT"] From 8a26f54d272d0e07c29bd0130a5d7039c6eb47d2 Mon Sep 17 00:00:00 2001 From: mytkom Date: Tue, 28 Jul 2026 16:23:02 +0200 Subject: [PATCH 07/12] add anchor prod CCDB timestamp getters --- internal/db/models/training_dataset.go | 12 +- internal/handler/training_dataset_handler.go | 4 +- internal/handler/training_task_handler.go | 4 +- internal/router.go | 5 +- internal/service/handler_errors.go | 7 + internal/service/monalisa_service.go | 64 +++++ internal/service/training_dataset_service.go | 21 +- .../training_dataset_timestamp_range.go | 107 ++++++++ internal/service/training_task_service.go | 149 ++--------- test/integration/common.go | 7 +- .../training_dataset_handler_test.go | 25 ++ .../integration/training_task_handler_test.go | 17 +- .../training_dataset_repository_test.go | 2 +- .../training_task_repository_test.go | 2 +- test/service/training_dataset_service_test.go | 16 +- test/service/training_task_service_test.go | 245 ++++++++---------- web/templates/training-datasets/new.html | 57 +--- 17 files changed, 387 insertions(+), 357 deletions(-) create mode 100644 internal/service/monalisa_service.go create mode 100644 internal/service/training_dataset_timestamp_range.go diff --git a/internal/db/models/training_dataset.go b/internal/db/models/training_dataset.go index d2c321e..fb95001 100644 --- a/internal/db/models/training_dataset.go +++ b/internal/db/models/training_dataset.go @@ -1,14 +1,18 @@ package models import ( + "time" + "github.com/mytkom/AliceTraINT/internal/jalien" "gorm.io/gorm" ) type TrainingDataset struct { gorm.Model - Name string `gorm:"type:varchar(255);not null;uniqueIndex"` - AODFiles []jalien.AODFile `gorm:"serializer:json"` - UserId uint - User User + Name string `gorm:"type:varchar(255);not null;uniqueIndex"` + AODFiles []jalien.AODFile `gorm:"serializer:json"` + AnchorProdTimestampFrom time.Time + AnchorProdTimestampTo time.Time + UserId uint + User User } diff --git a/internal/handler/training_dataset_handler.go b/internal/handler/training_dataset_handler.go index a66bf15..e170901 100644 --- a/internal/handler/training_dataset_handler.go +++ b/internal/handler/training_dataset_handler.go @@ -258,12 +258,12 @@ func (h *TrainingDatasetHandler) SelectRandomAods(w http.ResponseWriter, r *http } } -func InitTrainingDatasetRoutes(mux *http.ServeMux, env *environment.Env, jalien service.IJAliEnService, routeCache *utils.Cache) { +func InitTrainingDatasetRoutes(mux *http.ServeMux, env *environment.Env, jalien service.IJAliEnService, ccdbService service.ICCDBService, monalisaService service.IMonalisaService, routeCache *utils.Cache) { prefix := "training-datasets" tjh := &TrainingDatasetHandler{ Env: env, - Service: service.NewTrainingDatasetService(env.RepositoryContext, jalien), + Service: service.NewTrainingDatasetService(env.RepositoryContext, jalien, ccdbService, monalisaService), } cache := routeCache diff --git a/internal/handler/training_task_handler.go b/internal/handler/training_task_handler.go index f505dff..db5d270 100644 --- a/internal/handler/training_task_handler.go +++ b/internal/handler/training_task_handler.go @@ -174,10 +174,10 @@ func (h *TrainingTaskHandler) Create(w http.ResponseWriter, r *http.Request) { w.WriteHeader(http.StatusCreated) } -func InitTrainingTaskRoutes(mux *http.ServeMux, env *environment.Env, ccdbService service.ICCDBService, jalienService service.IJAliEnService, fileService service.IFileService, nnArch service.INNArchService) { +func InitTrainingTaskRoutes(mux *http.ServeMux, env *environment.Env, ccdbService service.ICCDBService, monalisaService service.IMonalisaService, fileService service.IFileService, nnArch service.INNArchService) { prefix := "training-tasks" - ttService := service.NewTrainingTaskService(env.RepositoryContext, ccdbService, jalienService, fileService, nnArch) + ttService := service.NewTrainingTaskService(env.RepositoryContext, ccdbService, monalisaService, fileService, nnArch) tjh := NewTrainingTaskHandler(env, ttService) authMw := middleware.NewAuthMw(env.IAuthService, true) diff --git a/internal/router.go b/internal/router.go index 82a0abe..4b36f97 100644 --- a/internal/router.go +++ b/internal/router.go @@ -29,6 +29,7 @@ func NewRouter(cfg *config.Config, repoContext *repository.RepositoryContext, au ccdbService := service.NewCCDBService(env) jalienCache := utils.NewCache(time.Duration(cfg.JalienCacheMinutes) * time.Minute) jalienService := service.NewJAliEnService(env, jalienCache) + monalisaService := service.NewMonalisaService(env) nnArch := service.NewNNArchService(cfg.NNArchPath) // local file storage fileService := service.NewLocalFileService(cfg.DataDirPath) @@ -44,8 +45,8 @@ func NewRouter(cfg *config.Config, repoContext *repository.RepositoryContext, au // handlers' routes handler.InitLandingRoutes(mux, env) handler.InitDocsRoutes(mux, env) - handler.InitTrainingDatasetRoutes(mux, env, jalienService, jalienCache) - handler.InitTrainingTaskRoutes(mux, env, ccdbService, jalienService, fileService, nnArch) + handler.InitTrainingDatasetRoutes(mux, env, jalienService, ccdbService, monalisaService, jalienCache) + handler.InitTrainingTaskRoutes(mux, env, ccdbService, monalisaService, fileService, nnArch) handler.InitTrainingMachineRoutes(mux, env, hasher) handler.InitQueueRoutes(mux, env, fileService, hasher) diff --git a/internal/service/handler_errors.go b/internal/service/handler_errors.go index d68257b..8ae8a24 100644 --- a/internal/service/handler_errors.go +++ b/internal/service/handler_errors.go @@ -18,6 +18,13 @@ func handleJAlienError(err error) error { return errJAlienUnreachable } +func handleMonalisaError(err error) error { + if os.IsTimeout(err) { + return errMonalisaUnreachable + } + return errInternalServerError +} + type ErrHandlerNotFound struct { Resource string } diff --git a/internal/service/monalisa_service.go b/internal/service/monalisa_service.go new file mode 100644 index 0000000..1256bb7 --- /dev/null +++ b/internal/service/monalisa_service.go @@ -0,0 +1,64 @@ +package service + +import ( + "sync" + + "github.com/mytkom/AliceTraINT/internal/environment" + "github.com/mytkom/AliceTraINT/internal/monalisa" + "github.com/stretchr/testify/mock" +) + +type IMonalisaService interface { + GetRunList() (*monalisa.HyperloopRunList, error) +} + +type MonalisaService struct { + client *monalisa.MonalisaClient + mu sync.RWMutex + cachedRunList *monalisa.HyperloopRunList +} + +func NewMonalisaService(env *environment.Env) *MonalisaService { + return &MonalisaService{ + client: monalisa.NewMonalisaClient(env.KeyPath, env.CertPath, env.MonalisaBaseURL), + } +} + +func (s *MonalisaService) GetRunList() (*monalisa.HyperloopRunList, error) { + s.mu.RLock() + if s.cachedRunList != nil { + runList := s.cachedRunList + s.mu.RUnlock() + return runList, nil + } + s.mu.RUnlock() + + runList, err := s.client.GetRunList() + if err != nil { + return nil, err + } + + s.mu.Lock() + s.cachedRunList = runList + s.mu.Unlock() + + return runList, nil +} + +type MockMonalisaService struct { + mock.Mock +} + +func NewMockMonalisaService() *MockMonalisaService { + return &MockMonalisaService{} +} + +func (s *MockMonalisaService) GetRunList() (*monalisa.HyperloopRunList, error) { + args := s.Called() + + if args.Error(1) != nil { + return nil, args.Error(1) + } + + return args.Get(0).(*monalisa.HyperloopRunList), args.Error(1) +} diff --git a/internal/service/training_dataset_service.go b/internal/service/training_dataset_service.go index 7457684..c2d4a45 100644 --- a/internal/service/training_dataset_service.go +++ b/internal/service/training_dataset_service.go @@ -23,17 +23,22 @@ type ITrainingDatasetService interface { type TrainingDatasetService struct { *repository.RepositoryContext - JAliEn IJAliEnService + JAliEn IJAliEnService + CCDBService ICCDBService + MonalisaService IMonalisaService } var errCCDBUnreachable = NewErrExternalServiceTimeout("CCDB") var errJAlienUnreachable = NewErrExternalServiceTimeout("JAlien") +var errMonalisaUnreachable = NewErrExternalServiceTimeout("Monalisa") var errDatasetNotFound = NewErrHandlerNotFound("TrainingDataset") -func NewTrainingDatasetService(repo *repository.RepositoryContext, jalien IJAliEnService) *TrainingDatasetService { +func NewTrainingDatasetService(repo *repository.RepositoryContext, jalien IJAliEnService, ccdbService ICCDBService, monalisaService IMonalisaService) *TrainingDatasetService { return &TrainingDatasetService{ RepositoryContext: repo, JAliEn: jalien, + CCDBService: ccdbService, + MonalisaService: monalisaService, } } @@ -78,7 +83,17 @@ func (s *TrainingDatasetService) Create(td *models.TrainingDataset) error { } } - err := s.TrainingDataset.Create(td) + timestampFrom, timestampTo, err := resolveDatasetTimestampRange(td.AODFiles, s.MonalisaService, s.CCDBService) + if err != nil { + if errors.Is(err, errCCDBUnreachable) || errors.Is(err, errInternalServerError) { + return err + } + return handleMonalisaError(err) + } + td.AnchorProdTimestampFrom = timestampFrom + td.AnchorProdTimestampTo = timestampTo + + err = s.TrainingDataset.Create(td) if err != nil { if errors.Is(err, gorm.ErrDuplicatedKey) { diff --git a/internal/service/training_dataset_timestamp_range.go b/internal/service/training_dataset_timestamp_range.go new file mode 100644 index 0000000..bf455d9 --- /dev/null +++ b/internal/service/training_dataset_timestamp_range.go @@ -0,0 +1,107 @@ +package service + +import ( + "errors" + "time" + + "github.com/mytkom/AliceTraINT/internal/ccdb" + "github.com/mytkom/AliceTraINT/internal/jalien" +) + +func resolveDatasetTimestampRange(aodFiles []jalien.AODFile, monalisaService IMonalisaService, ccdbService ICCDBService) (time.Time, time.Time, error) { + if len(aodFiles) == 0 { + return time.Time{}, time.Time{}, errors.New("unexpected behaviour: empty training dataset") + } + + distinctPeriods := make(map[string]struct{}, len(aodFiles)) + for _, aod := range aodFiles { + if aod.LHCPeriod != "" { + distinctPeriods[aod.LHCPeriod] = struct{}{} + } + } + + if len(distinctPeriods) == 0 { + return time.Time{}, time.Time{}, errors.New("unexpected behaviour: empty training dataset periods") + } + + runList, err := monalisaService.GetRunList() + if err != nil { + return time.Time{}, time.Time{}, err + } + + globalMinRun, globalMaxRun, err := findGlobalRunRange(runList.TagToRuns, distinctPeriods) + if err != nil { + return time.Time{}, time.Time{}, err + } + + firstRunInfo, lastRunInfo, err := getRunInfoRangeFromCCDB(ccdbService, globalMinRun, globalMaxRun) + if err != nil { + return time.Time{}, time.Time{}, err + } + + return time.UnixMilli(int64(firstRunInfo.SOR)).UTC(), time.UnixMilli(int64(lastRunInfo.EOR)).UTC(), nil +} + +func findGlobalRunRange(runMap map[string][]uint64, periods map[string]struct{}) (uint64, uint64, error) { + var globalMinRun uint64 + var globalMaxRun uint64 + initialized := false + + for period := range periods { + runs, ok := runMap[period] + if !ok || len(runs) == 0 { + continue + } + + for _, run := range runs { + if !initialized || run < globalMinRun { + globalMinRun = run + } + if !initialized || run > globalMaxRun { + globalMaxRun = run + } + initialized = true + } + } + + if !initialized { + return 0, 0, errors.New("no runs found for dataset periods in Monalisa run list") + } + + return globalMinRun, globalMaxRun, nil +} + +func getRunInfoRangeFromCCDB(ccdbService ICCDBService, smallestRun, greatestRun uint64) (*ccdb.RunInformation, *ccdb.RunInformation, error) { + firstRunInfo, err := ccdbService.GetRunInformation(smallestRun) + if err != nil { + return nil, nil, handleCCDBError(err) + } + + lastRunInfo, err := ccdbService.GetRunInformation(greatestRun) + if err != nil { + return nil, nil, handleCCDBError(err) + } + + return firstRunInfo, lastRunInfo, nil +} + +func getDatasetAODRunRange(aodFiles []jalien.AODFile) (uint64, uint64, error) { + var minRun, maxRun uint64 + initialized := false + + for _, aod := range aodFiles { + if !initialized || aod.RunNumber < minRun { + minRun = aod.RunNumber + } + if !initialized || aod.RunNumber > maxRun { + maxRun = aod.RunNumber + } + initialized = true + } + + if !initialized { + return 0, 0, errors.New("unexpected behaviour: empty training dataset") + } + + return minRun, maxRun, nil +} diff --git a/internal/service/training_task_service.go b/internal/service/training_task_service.go index 5017b65..f27b223 100644 --- a/internal/service/training_task_service.go +++ b/internal/service/training_task_service.go @@ -4,14 +4,9 @@ import ( "errors" "fmt" "log" - "regexp" - "slices" - "strconv" - "github.com/mytkom/AliceTraINT/internal/ccdb" "github.com/mytkom/AliceTraINT/internal/db/models" "github.com/mytkom/AliceTraINT/internal/db/repository" - "github.com/mytkom/AliceTraINT/internal/jalien" "gorm.io/gorm" ) @@ -37,21 +32,19 @@ type ITrainingTaskService interface { type TrainingTaskService struct { *repository.RepositoryContext - CCDBService ICCDBService - JAliEnService IJAliEnService - FileService IFileService - NNArch INNArchService - PeriodRegex *regexp.Regexp + CCDBService ICCDBService + MonalisaService IMonalisaService + FileService IFileService + NNArch INNArchService } -func NewTrainingTaskService(repo *repository.RepositoryContext, ccdbService ICCDBService, jalienService IJAliEnService, fileService IFileService, nnArch INNArchService) *TrainingTaskService { +func NewTrainingTaskService(repo *repository.RepositoryContext, ccdbService ICCDBService, monalisaService IMonalisaService, fileService IFileService, nnArch INNArchService) *TrainingTaskService { return &TrainingTaskService{ RepositoryContext: repo, CCDBService: ccdbService, - JAliEnService: jalienService, + MonalisaService: monalisaService, FileService: fileService, NNArch: nnArch, - PeriodRegex: regexp.MustCompile(`(/alice/sim/\d{4}/LHC[a-z0-9A-Z\_].+(/\d+)?)/\d+/AOD/\d+`), } } @@ -166,43 +159,11 @@ func (s *TrainingTaskService) UploadOnnxResults(id uint) error { } } - lhcPeriods, err := s.getLHCPeriods(trainingTask) + minSOR, maxEOR, err := s.getUploadTimestampRange(trainingTask) if err != nil { return err } - var minSOR, maxEOR uint64 - initialized := false - - for i, period := range lhcPeriods { - log.Printf("%d: Name=\"%s\" DirPath=\"%s\"", i, period.Name, period.DirPath) - - dirContents, err := s.JAliEnService.ListAndParseDirectory(period.DirPath) - if err != nil { - return err - } - - smallestRun, greatestRun, err := s.findRunNumberRange(dirContents.Subdirs) - if err != nil { - return err - } - - firstRunInfo, lastRunInfo, err := s.getRunInfoRange(smallestRun, greatestRun) - if err != nil { - return err - } - - if !initialized || firstRunInfo.SOR < minSOR { - minSOR = firstRunInfo.SOR - } - - if !initialized || lastRunInfo.EOR > maxEOR { - maxEOR = lastRunInfo.EOR - } - - initialized = true - } - mappedOnnxFiles, err := s.filterOnnxFiles(trainingTask.ID) if err != nil { return err @@ -222,98 +183,32 @@ func (s *TrainingTaskService) UploadOnnxResults(id uint) error { return nil } -type lhcPeriod struct { - Name string - DirPath string -} - -func (s *TrainingTaskService) periodPathFromAODPath(aodPath string) (string, error) { - matches := s.PeriodRegex.FindStringSubmatch(aodPath) - - if len(matches) != 3 { - return "", errors.New("unexpected AOD path format, cannot correctly match") - } - - return matches[1], nil -} - -func (s *TrainingTaskService) getLHCPeriods(task *models.TrainingTask) ([]lhcPeriod, error) { - var periods []lhcPeriod - initialized := false - - for _, aod := range task.TrainingDataset.AODFiles { - if !slices.ContainsFunc(periods, func(p lhcPeriod) bool { - return p.Name == aod.LHCPeriod - }) { - periodPath, err := s.periodPathFromAODPath(aod.Path) - if err != nil { - return nil, err - } - - periods = append(periods, lhcPeriod{ - Name: aod.LHCPeriod, - DirPath: periodPath, - }) - } - initialized = true +func (s *TrainingTaskService) getUploadTimestampRange(trainingTask *models.TrainingTask) (uint64, uint64, error) { + if !trainingTask.TrainingDataset.AnchorProdTimestampFrom.IsZero() && + !trainingTask.TrainingDataset.AnchorProdTimestampTo.IsZero() { + return uint64(trainingTask.TrainingDataset.AnchorProdTimestampFrom.UnixMilli()), + uint64(trainingTask.TrainingDataset.AnchorProdTimestampTo.UnixMilli()), + nil } - if !initialized { - return nil, errors.New("unexpected behaviour: empty training dataset") + timestampFrom, timestampTo, err := resolveDatasetTimestampRange(trainingTask.TrainingDataset.AODFiles, s.MonalisaService, s.CCDBService) + if err == nil { + return uint64(timestampFrom.UnixMilli()), uint64(timestampTo.UnixMilli()), nil } - return periods, nil -} - -func (s *TrainingTaskService) findRunNumberRange(subdirs []jalien.Dir) (uint64, uint64, error) { - var smallestRun, greatestRun uint64 - initialized := false + log.Printf("timestamp resolver failed, using dataset-run fallback: %v", err) - for _, dir := range subdirs { - log.Printf("Dir: %+v", dir) - - runNumber, err := strconv.ParseUint(dir.Name, 10, 64) - if err != nil { - log.Println(err.Error()) - continue - } - log.Printf("RunNumber: %+v", runNumber) - - if !initialized || runNumber < smallestRun { - smallestRun = runNumber - } - if !initialized || runNumber > greatestRun { - greatestRun = runNumber - } - initialized = true - } - - log.Printf("SmallestRun: %+v", smallestRun) - log.Printf("GreatestRun: %+v", greatestRun) - log.Printf("Initialized: %+v", initialized) - - if !initialized { - return 0, 0, errors.New("unexpected behaviour: empty training dataset") - } - - return smallestRun, greatestRun, nil -} - -func (s *TrainingTaskService) getRunInfoRange(smallestRun, greatestRun uint64) (*ccdb.RunInformation, *ccdb.RunInformation, error) { - firstRunInfo, err := s.CCDBService.GetRunInformation(smallestRun) + minRun, maxRun, err := getDatasetAODRunRange(trainingTask.TrainingDataset.AODFiles) if err != nil { - return nil, nil, handleCCDBError(err) + return 0, 0, err } - lastRunInfo, err := s.CCDBService.GetRunInformation(greatestRun) + firstRunInfo, lastRunInfo, err := getRunInfoRangeFromCCDB(s.CCDBService, minRun, maxRun) if err != nil { - return nil, nil, handleCCDBError(err) + return 0, 0, err } - log.Printf("From run %d, SOR %d", firstRunInfo.RunNumber, firstRunInfo.SOR) - log.Printf("to run %d, EOR %d", lastRunInfo.RunNumber, lastRunInfo.EOR) - - return firstRunInfo, lastRunInfo, nil + return firstRunInfo.SOR, lastRunInfo.EOR, nil } func (s *TrainingTaskService) filterOnnxFiles(ttId uint) (map[string]*models.TrainingTaskResult, error) { diff --git a/test/integration/common.go b/test/integration/common.go index 838d909..93461cc 100644 --- a/test/integration/common.go +++ b/test/integration/common.go @@ -29,6 +29,7 @@ func init() { type MockedServices struct { CCDB *service.MockCCDBService JAliEn *service.MockJAliEnService + Monalisa *service.MockMonalisaService FileService *service.MockFileService Auth *auth.AuthServiceMock NNArch *service.NNArchServiceInMemory @@ -54,6 +55,7 @@ func mockRouter(db *gorm.DB, cfg *config.Config) *IntegrationTestUtils { hasher := service.NewMockHasher() ccdbService := service.NewMockCCDBService() jalienService := service.NewMockJAliEnService() + monalisaService := service.NewMockMonalisaService() nnArch := service.NewNNArchServiceInMemory(&service.NNFieldConfigs{ "fieldName": service.NNConfigField{ FullName: "Full field name", @@ -73,8 +75,8 @@ func mockRouter(db *gorm.DB, cfg *config.Config) *IntegrationTestUtils { // handlers' routes handler.InitLandingRoutes(mux, env) - handler.InitTrainingDatasetRoutes(mux, env, jalienService, nil) - handler.InitTrainingTaskRoutes(mux, env, ccdbService, jalienService, fileService, nnArch) + handler.InitTrainingDatasetRoutes(mux, env, jalienService, ccdbService, monalisaService, nil) + handler.InitTrainingTaskRoutes(mux, env, ccdbService, monalisaService, fileService, nnArch) handler.InitTrainingMachineRoutes(mux, env, hasher) handler.InitQueueRoutes(mux, env, fileService, hasher) @@ -84,6 +86,7 @@ func mockRouter(db *gorm.DB, cfg *config.Config) *IntegrationTestUtils { MockedServices: &MockedServices{ CCDB: ccdbService, JAliEn: jalienService, + Monalisa: monalisaService, FileService: fileService, Auth: auth, NNArch: nnArch, diff --git a/test/integration/training_dataset_handler_test.go b/test/integration/training_dataset_handler_test.go index 46bca4b..d3e48c3 100644 --- a/test/integration/training_dataset_handler_test.go +++ b/test/integration/training_dataset_handler_test.go @@ -7,9 +7,12 @@ import ( "fmt" "net/http" "testing" + "time" + "github.com/mytkom/AliceTraINT/internal/ccdb" "github.com/mytkom/AliceTraINT/internal/db/models" "github.com/mytkom/AliceTraINT/internal/jalien" + "github.com/mytkom/AliceTraINT/internal/monalisa" "github.com/stretchr/testify/assert" ) @@ -165,6 +168,24 @@ type createTrainingDatasetPayload struct { UserId uint } +func mockDatasetTimestampRange(ut *IntegrationTestUtils, period string, minRun, maxRun uint64) { + ut.Monalisa.On("GetRunList").Return(&monalisa.HyperloopRunList{ + TagToRuns: map[string][]uint64{ + period: {minRun, maxRun}, + }, + }, nil).Once() + ut.CCDB.On("GetRunInformation", minRun).Return(&ccdb.RunInformation{ + RunNumber: minRun, + SOR: uint64(time.Now().Add(-time.Hour).UnixMilli()), + EOR: uint64(time.Now().Add(-30 * time.Minute).UnixMilli()), + }, nil).Once() + ut.CCDB.On("GetRunInformation", maxRun).Return(&ccdb.RunInformation{ + RunNumber: maxRun, + SOR: uint64(time.Now().Add(-20 * time.Minute).UnixMilli()), + EOR: uint64(time.Now().UnixMilli()), + }, nil).Once() +} + func TestTrainingDatasetHandler_Create_Failure(t *testing.T) { ut, cleanup := setupIntegrationTest(t) defer cleanup() @@ -201,6 +222,8 @@ func TestTrainingDatasetHandler_Create_Failure(t *testing.T) { body, err := json.Marshal(tdSameName) assert.NoError(t, err) + mockDatasetTimestampRange(ut, "LHC24b1b", 567454, 567456) + req, err := http.NewRequest("POST", "/training-datasets", bytes.NewReader(body)) assert.NoError(t, err) HTMXReq(req) @@ -236,6 +259,8 @@ func TestTrainingDatasetHandler_Create(t *testing.T) { body, err := json.Marshal(trainingDataset) assert.NoError(t, err) + mockDatasetTimestampRange(ut, "LHC24b1b", 567454, 567454) + req, err := http.NewRequest("POST", "/training-datasets", bytes.NewReader(body)) assert.NoError(t, err) HTMXReq(req) diff --git a/test/integration/training_task_handler_test.go b/test/integration/training_task_handler_test.go index 4662f0d..0feb2bf 100644 --- a/test/integration/training_task_handler_test.go +++ b/test/integration/training_task_handler_test.go @@ -293,21 +293,14 @@ func prepareUploadToCCDB(t *testing.T, ut *IntegrationTestUtils, user *models.Us assert.NoError(t, ut.TrainingTask.Create(trainingTask)) now := uint64(time.Now().UTC().UnixMilli()) - ut.MockedServices.JAliEn.On("ListAndParseDirectory", "/alice/sim/2024/LHC24b1b/0").Return(&jalien.DirectoryContents{ - Subdirs: []jalien.Dir{ - {Name: "560000", Path: "/alice/sim/2024/LHC24b1b/0/560000"}, - {Name: "567454", Path: "/alice/sim/2024/LHC24b1b/0/567454"}, - {Name: "567458", Path: "/alice/sim/2024/LHC24b1b/0/567458"}, - {Name: "570000", Path: "/alice/sim/2024/LHC24b1b/0/570000"}, - }, - }, nil) - ut.MockedServices.CCDB.On("GetRunInformation", uint64(560000)).Return(&ccdb.RunInformation{ - RunNumber: 560000, + ut.MockedServices.Monalisa.On("GetRunList").Return(nil, errors.New("monalisa down")) + ut.MockedServices.CCDB.On("GetRunInformation", uint64(567454)).Return(&ccdb.RunInformation{ + RunNumber: 567454, SOR: now - 10000, EOR: now, }, nil) - ut.MockedServices.CCDB.On("GetRunInformation", uint64(570000)).Return(&ccdb.RunInformation{ - RunNumber: 570000, + ut.MockedServices.CCDB.On("GetRunInformation", uint64(567458)).Return(&ccdb.RunInformation{ + RunNumber: 567458, SOR: now, EOR: now + 10000, }, nil) diff --git a/test/repository/training_dataset_repository_test.go b/test/repository/training_dataset_repository_test.go index c75778b..a8edb13 100644 --- a/test/repository/training_dataset_repository_test.go +++ b/test/repository/training_dataset_repository_test.go @@ -40,7 +40,7 @@ func TestTrainingDatasetRepository_Create(t *testing.T) { mock.ExpectBegin() mock.ExpectQuery(`INSERT INTO "training_datasets" (.+) RETURNING "id"`). - WithArgs(AnyTime(), AnyTime(), AnyTime(), trainingDataset.Name, marshalAODFiles(t, trainingDataset), 1). + WithArgs(AnyTime(), AnyTime(), AnyTime(), trainingDataset.Name, marshalAODFiles(t, trainingDataset), AnyTime(), AnyTime(), 1). WillReturnRows(sqlmock.NewRows([]string{"id"}).AddRow(1)) mock.ExpectCommit() diff --git a/test/repository/training_task_repository_test.go b/test/repository/training_task_repository_test.go index ae4a8f6..6776241 100644 --- a/test/repository/training_task_repository_test.go +++ b/test/repository/training_task_repository_test.go @@ -48,7 +48,7 @@ func TestTrainingTaskRepository_Create(t *testing.T) { mock.ExpectBegin() mock.ExpectQuery(`INSERT INTO "training_datasets" (.+) RETURNING "id"`). - WithArgs(AnyTime(), AnyTime(), AnyTime(), trainingDataset.Name, marshalAODFiles(t, trainingDataset), 1). + WithArgs(AnyTime(), AnyTime(), AnyTime(), trainingDataset.Name, marshalAODFiles(t, trainingDataset), AnyTime(), AnyTime(), 1). WillReturnRows(sqlmock.NewRows([]string{"id"}).AddRow(1)) mock.ExpectQuery(`INSERT INTO "training_tasks" (.+) RETURNING "id"`). WithArgs(AnyTime(), AnyTime(), AnyTime(), trainingTask.Name, trainingTask.Status, 1, 1, nil, marshalTrainingTaskConfig(t, trainingTask)). diff --git a/test/service/training_dataset_service_test.go b/test/service/training_dataset_service_test.go index e768be8..e3d9705 100644 --- a/test/service/training_dataset_service_test.go +++ b/test/service/training_dataset_service_test.go @@ -11,18 +11,24 @@ import ( ) type trainingDatasetServiceTestUtils struct { - TDRepo *repository.MockTrainingDatasetRepository - JAliEnService *service.MockJAliEnService + TDRepo *repository.MockTrainingDatasetRepository + JAliEnService *service.MockJAliEnService + CCDBService *service.MockCCDBService + MonalisaService *service.MockMonalisaService } func newTrainingDatasetService() (*service.TrainingDatasetService, *trainingDatasetServiceTestUtils) { tdRepo := repository.NewMockTrainingDatasetRepository() jalienService := service.NewMockJAliEnService() + ccdbService := service.NewMockCCDBService() + monalisaService := service.NewMockMonalisaService() return service.NewTrainingDatasetService(&repository.RepositoryContext{ TrainingDataset: tdRepo, - }, jalienService), &trainingDatasetServiceTestUtils{ - TDRepo: tdRepo, - JAliEnService: jalienService, + }, jalienService, ccdbService, monalisaService), &trainingDatasetServiceTestUtils{ + TDRepo: tdRepo, + JAliEnService: jalienService, + CCDBService: ccdbService, + MonalisaService: monalisaService, } } diff --git a/test/service/training_task_service_test.go b/test/service/training_task_service_test.go index 43bf4da..320c825 100644 --- a/test/service/training_task_service_test.go +++ b/test/service/training_task_service_test.go @@ -11,6 +11,7 @@ import ( "github.com/mytkom/AliceTraINT/internal/db/models" "github.com/mytkom/AliceTraINT/internal/db/repository" "github.com/mytkom/AliceTraINT/internal/jalien" + "github.com/mytkom/AliceTraINT/internal/monalisa" "github.com/mytkom/AliceTraINT/internal/service" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/mock" @@ -18,21 +19,21 @@ import ( ) type trainingTaskServiceTestUtils struct { - TTRepo *repository.MockTrainingTaskRepository - TDRepo *repository.MockTrainingDatasetRepository - TTRRepo *repository.MockTrainingTaskResultRepository - CCDBService *service.MockCCDBService - JAliEnService *service.MockJAliEnService - FileService *service.MockFileService - NNArch *service.NNArchServiceInMemory + TTRepo *repository.MockTrainingTaskRepository + TDRepo *repository.MockTrainingDatasetRepository + TTRRepo *repository.MockTrainingTaskResultRepository + CCDBService *service.MockCCDBService + MonalisaService *service.MockMonalisaService + FileService *service.MockFileService + NNArch *service.NNArchServiceInMemory } func newTrainingTaskService() (*service.TrainingTaskService, *trainingTaskServiceTestUtils) { ttRepo := repository.NewMockTrainingTaskRepository() tdRepo := repository.NewMockTrainingDatasetRepository() ttrRepo := repository.NewMockTrainingTaskResultRepository() - jalienService := service.NewMockJAliEnService() ccdbService := service.NewMockCCDBService() + monalisaService := service.NewMockMonalisaService() fileService := service.NewMockFileService() nnArch := service.NewNNArchServiceInMemory(&service.NNFieldConfigs{ "fieldName": service.NNConfigField{ @@ -54,14 +55,14 @@ func newTrainingTaskService() (*service.TrainingTaskService, *trainingTaskServic TrainingTask: ttRepo, TrainingDataset: tdRepo, TrainingTaskResult: ttrRepo, - }, ccdbService, jalienService, fileService, nnArch), &trainingTaskServiceTestUtils{ - TTRepo: ttRepo, - TDRepo: tdRepo, - TTRRepo: ttrRepo, - CCDBService: ccdbService, - JAliEnService: jalienService, - FileService: fileService, - NNArch: nnArch, + }, ccdbService, monalisaService, fileService, nnArch), &trainingTaskServiceTestUtils{ + TTRepo: ttRepo, + TDRepo: tdRepo, + TTRRepo: ttrRepo, + CCDBService: ccdbService, + MonalisaService: monalisaService, + FileService: fileService, + NNArch: nnArch, } } @@ -396,27 +397,17 @@ func TestTrainingTaskService_UploadToCCDB_OnePeriod(t *testing.T) { ut.TTRepo.On("GetByID", ttId).Return(&tt, nil) ut.TTRepo.On("Update", &tt).Return(nil) ut.TTRRepo.On("GetByType", ttId, models.Onnx).Return(onnxFiles, nil) + ut.MonalisaService.On("GetRunList").Return(nil, errors.New("monalisa down")).Once() file := &mockReadCloser{} ut.FileService.On("OpenFile", "./local_file_temp.onnx").Return(file, func(r io.ReadCloser) { r.Close() }, nil) - ut.JAliEnService.On("ListAndParseDirectory", "/alice/sim/2024/LHC24f3/0").Return(&jalien.DirectoryContents{ - Subdirs: []jalien.Dir{ - {Name: "321000", Path: "/alice/sim/2024/LHC24f3/0/321000"}, - {Name: "321100", Path: "/alice/sim/2024/LHC24f3/0/321100"}, - {Name: "321321", Path: "/alice/sim/2024/LHC24f3/0/321321"}, - {Name: "321326", Path: "/alice/sim/2024/LHC24f3/0/321326"}, - {Name: "321338", Path: "/alice/sim/2024/LHC24f3/0/321338"}, - {Name: "321400", Path: "/alice/sim/2024/LHC24f3/0/321400"}, - {Name: "321500", Path: "/alice/sim/2024/LHC24f3/0/321500"}, - }, - }, nil) now := time.Now().UnixMilli() - ut.CCDBService.On("GetRunInformation", uint64(321000)).Return(&ccdb.RunInformation{ - RunNumber: 321000, + ut.CCDBService.On("GetRunInformation", uint64(321321)).Return(&ccdb.RunInformation{ + RunNumber: 321321, SOR: uint64(now - 10000), EOR: uint64(now - 9000), }, nil) - ut.CCDBService.On("GetRunInformation", uint64(321500)).Return(&ccdb.RunInformation{ - RunNumber: 321500, + ut.CCDBService.On("GetRunInformation", uint64(321338)).Return(&ccdb.RunInformation{ + RunNumber: 321338, SOR: uint64(now + 7000), EOR: uint64(now + 10000), }, nil) @@ -432,9 +423,8 @@ func TestTrainingTaskService_UploadToCCDB_OnePeriod(t *testing.T) { assert.Equal(t, models.Uploaded, tt.Status) ut.TTRRepo.AssertCalled(t, "GetByType", ttId, models.Onnx) ut.FileService.AssertCalled(t, "OpenFile", "./local_file_temp.onnx") - ut.JAliEnService.AssertCalled(t, "ListAndParseDirectory", "/alice/sim/2024/LHC24f3/0") - ut.CCDBService.AssertCalled(t, "GetRunInformation", uint64(321000)) - ut.CCDBService.AssertCalled(t, "GetRunInformation", uint64(321500)) + ut.CCDBService.AssertCalled(t, "GetRunInformation", uint64(321321)) + ut.CCDBService.AssertCalled(t, "GetRunInformation", uint64(321338)) ut.CCDBService.AssertCalled(t, "UploadFile", uint64(now-10000), uint64(now+10000), "uploaded_file.onnx", file) } @@ -478,75 +468,23 @@ func TestTrainingTaskService_UploadToCCDB_ManyPeriods(t *testing.T) { ut.TTRepo.On("GetByID", ttId).Return(&tt, nil) ut.TTRepo.On("Update", &tt).Return(nil) ut.TTRRepo.On("GetByType", ttId, models.Onnx).Return(onnxFiles, nil) + ut.MonalisaService.On("GetRunList").Return(nil, errors.New("monalisa down")).Once() file := &mockReadCloser{} ut.FileService.On("OpenFile", "./local_file_temp.onnx").Return(file, func(r io.ReadCloser) { r.Close() }, nil) - // Mock JAliEn and CCDB with timestamp information now := time.Now().UnixMilli() - // LHC24b1b - ut.JAliEnService.On("ListAndParseDirectory", "/alice/sim/2024/LHC24b1b/0").Return(&jalien.DirectoryContents{ - Subdirs: []jalien.Dir{ - {Name: "319900", Path: "/alice/sim/2024/LHC24b1b/0/319900"}, // min run number LHC24b1b - {Name: "320000", Path: "/alice/sim/2024/LHC24b1b/0/320000"}, - {Name: "320100", Path: "/alice/sim/2024/LHC24b1b/0/320100"}, - {Name: "320200", Path: "/alice/sim/2024/LHC24b1b/0/320200"}, - {Name: "320300", Path: "/alice/sim/2024/LHC24b1b/0/320300"}, // max run number LHC24b1b - }, - }, nil) - ut.CCDBService.On("GetRunInformation", uint64(319900)).Return(&ccdb.RunInformation{ - RunNumber: 319900, + ut.CCDBService.On("GetRunInformation", uint64(320000)).Return(&ccdb.RunInformation{ + RunNumber: 320000, SOR: uint64(now - 20000), EOR: uint64(now - 19000), }, nil) - ut.CCDBService.On("GetRunInformation", uint64(320300)).Return(&ccdb.RunInformation{ - RunNumber: 320300, - SOR: uint64(now - 12000), - EOR: uint64(now - 11000), - }, nil) - // LHC24f3 - ut.JAliEnService.On("ListAndParseDirectory", "/alice/sim/2024/LHC24f3/0").Return(&jalien.DirectoryContents{ - Subdirs: []jalien.Dir{ - {Name: "321000", Path: "/alice/sim/2024/LHC24f3/0/321000"}, // min run number LHC24f3 - {Name: "321321", Path: "/alice/sim/2024/LHC24f3/0/321321"}, - {Name: "321326", Path: "/alice/sim/2024/LHC24f3/0/321326"}, - {Name: "321338", Path: "/alice/sim/2024/LHC24f3/0/321338"}, - {Name: "321500", Path: "/alice/sim/2024/LHC24f3/0/321500"}, // max run number LHC24f3 - }, - }, nil) - ut.CCDBService.On("GetRunInformation", uint64(321000)).Return(&ccdb.RunInformation{ - RunNumber: 321000, - SOR: uint64(now - 10000), - EOR: uint64(now - 9000), - }, nil) - ut.CCDBService.On("GetRunInformation", uint64(321500)).Return(&ccdb.RunInformation{ - RunNumber: 321500, - SOR: uint64(now - 2000), - EOR: uint64(now - 1000), - }, nil) - // LHC24c1 - // simulate 2023 datasets (without numbered subdir in period dir) - ut.JAliEnService.On("ListAndParseDirectory", "/alice/sim/2024/LHC24c1").Return(&jalien.DirectoryContents{ - Subdirs: []jalien.Dir{ - {Name: "321900", Path: "/alice/sim/2024/LHC24c1/321900"}, // min run number LHC24c1 - {Name: "322000", Path: "/alice/sim/2024/LHC24c1/322000"}, - {Name: "322100", Path: "/alice/sim/2024/LHC24c1/322100"}, - {Name: "322200", Path: "/alice/sim/2024/LHC24c1/322200"}, - {Name: "322300", Path: "/alice/sim/2024/LHC24c1/322300"}, // max run number LHC24c1 - }, - }, nil) - ut.CCDBService.On("GetRunInformation", uint64(321900)).Return(&ccdb.RunInformation{ - RunNumber: 321900, - SOR: uint64(now), - EOR: uint64(now + 1000), - }, nil) - ut.CCDBService.On("GetRunInformation", uint64(322300)).Return(&ccdb.RunInformation{ - RunNumber: 322300, + ut.CCDBService.On("GetRunInformation", uint64(322200)).Return(&ccdb.RunInformation{ + RunNumber: 322200, SOR: uint64(now + 9000), EOR: uint64(now + 10000), }, nil) - // TODO b1b and c1 ut.CCDBService.On("UploadFile", uint64(now-20000), uint64(now+10000), "uploaded_file.onnx", file).Return(nil) // Act @@ -559,19 +497,68 @@ func TestTrainingTaskService_UploadToCCDB_ManyPeriods(t *testing.T) { assert.Equal(t, models.Uploaded, tt.Status) ut.TTRRepo.AssertCalled(t, "GetByType", ttId, models.Onnx) ut.FileService.AssertCalled(t, "OpenFile", "./local_file_temp.onnx") - // LHC24b1b info - ut.JAliEnService.AssertCalled(t, "ListAndParseDirectory", "/alice/sim/2024/LHC24b1b/0") - ut.CCDBService.AssertCalled(t, "GetRunInformation", uint64(319900)) - ut.CCDBService.AssertCalled(t, "GetRunInformation", uint64(320300)) - // LHC24f3 info - ut.JAliEnService.AssertCalled(t, "ListAndParseDirectory", "/alice/sim/2024/LHC24f3/0") - ut.CCDBService.AssertCalled(t, "GetRunInformation", uint64(321000)) - ut.CCDBService.AssertCalled(t, "GetRunInformation", uint64(321500)) - // LHC24c1 info - ut.JAliEnService.AssertCalled(t, "ListAndParseDirectory", "/alice/sim/2024/LHC24c1") - ut.CCDBService.AssertCalled(t, "GetRunInformation", uint64(321900)) - ut.CCDBService.AssertCalled(t, "GetRunInformation", uint64(322300)) - // Upload + ut.CCDBService.AssertCalled(t, "GetRunInformation", uint64(320000)) + ut.CCDBService.AssertCalled(t, "GetRunInformation", uint64(322200)) + ut.CCDBService.AssertCalled(t, "UploadFile", uint64(now-20000), uint64(now+10000), "uploaded_file.onnx", file) +} + +func TestTrainingTaskService_UploadToCCDB_UsesMonalisaFirst(t *testing.T) { + ttService, ut := newTrainingTaskService() + userId := uint(1) + tdId := uint(1) + ttId := uint(1) + tt := models.TrainingTask{ + Model: gorm.Model{ID: ttId}, + Name: "task2", + UserId: userId, + Status: models.Completed, + TrainingDatasetId: tdId, + TrainingDataset: models.TrainingDataset{ + Model: gorm.Model{ID: tdId}, + UserId: userId, + AODFiles: []jalien.AODFile{ + {Name: "AO2D.root", Path: "/alice/sim/2024/LHC24b1b/0/320000/AOD/002", RunNumber: 320000, LHCPeriod: "LHC24b1b", AODNumber: 2}, + {Name: "AO2D.root", Path: "/alice/sim/2024/LHC24f3/0/321338/AOD/002", RunNumber: 321338, LHCPeriod: "LHC24f3", AODNumber: 2}, + }, + }, + Configuration: "", + } + onnxFiles := []models.TrainingTaskResult{ + {Name: "local_file.onnx", Type: models.Onnx, Description: "example", FileId: 1, File: models.File{ + Name: "local_file_temp.onnx", Path: "./local_file_temp.onnx", Size: 12312, + }, TrainingTaskId: ttId}, + } + + ut.TTRepo.On("GetByID", ttId).Return(&tt, nil) + ut.TTRepo.On("Update", &tt).Return(nil) + ut.TTRRepo.On("GetByType", ttId, models.Onnx).Return(onnxFiles, nil) + ut.MonalisaService.On("GetRunList").Return(&monalisa.HyperloopRunList{ + TagToRuns: map[string][]uint64{ + "LHC24b1b": {319900, 320300}, + "LHC24f3": {321000, 321500}, + }, + }, nil).Once() + + file := &mockReadCloser{} + ut.FileService.On("OpenFile", "./local_file_temp.onnx").Return(file, func(r io.ReadCloser) { r.Close() }, nil) + + now := time.Now().UnixMilli() + ut.CCDBService.On("GetRunInformation", uint64(319900)).Return(&ccdb.RunInformation{ + RunNumber: 319900, + SOR: uint64(now - 20000), + EOR: uint64(now - 19000), + }, nil) + ut.CCDBService.On("GetRunInformation", uint64(321500)).Return(&ccdb.RunInformation{ + RunNumber: 321500, + SOR: uint64(now + 9000), + EOR: uint64(now + 10000), + }, nil) + ut.CCDBService.On("UploadFile", uint64(now-20000), uint64(now+10000), "uploaded_file.onnx", file).Return(nil) + + err := ttService.UploadOnnxResults(ttId) + + assert.NoError(t, err) + ut.MonalisaService.AssertCalled(t, "GetRunList") ut.CCDBService.AssertCalled(t, "UploadFile", uint64(now-20000), uint64(now+10000), "uploaded_file.onnx", file) } @@ -606,27 +593,17 @@ func TestTrainingTaskService_UploadToCCDB_MissingExpectedFile(t *testing.T) { ut.TTRepo.On("GetByID", ttId).Return(&tt, nil) ut.TTRepo.On("Update", &tt).Return(nil) ut.TTRRepo.On("GetByType", ttId, models.Onnx).Return(onnxFiles, nil) + ut.MonalisaService.On("GetRunList").Return(nil, errors.New("monalisa down")).Once() file := &mockReadCloser{} ut.FileService.On("OpenFile", "./local_file_temp.onnx").Return(file, func(r io.ReadCloser) { r.Close() }, nil) - ut.JAliEnService.On("ListAndParseDirectory", "/alice/sim/2024/LHC24f3/0").Return(&jalien.DirectoryContents{ - Subdirs: []jalien.Dir{ - {Name: "321000", Path: "/alice/sim/2024/LHC24f3/0/321000"}, - {Name: "321100", Path: "/alice/sim/2024/LHC24f3/0/321100"}, - {Name: "321321", Path: "/alice/sim/2024/LHC24f3/0/321321"}, - {Name: "321326", Path: "/alice/sim/2024/LHC24f3/0/321326"}, - {Name: "321338", Path: "/alice/sim/2024/LHC24f3/0/321338"}, - {Name: "321400", Path: "/alice/sim/2024/LHC24f3/0/321400"}, - {Name: "321500", Path: "/alice/sim/2024/LHC24f3/0/321500"}, - }, - }, nil) now := time.Now().UnixMilli() - ut.CCDBService.On("GetRunInformation", uint64(321000)).Return(&ccdb.RunInformation{ - RunNumber: 321000, + ut.CCDBService.On("GetRunInformation", uint64(321321)).Return(&ccdb.RunInformation{ + RunNumber: 321321, SOR: uint64(now - 10000), EOR: uint64(now - 9000), }, nil) - ut.CCDBService.On("GetRunInformation", uint64(321500)).Return(&ccdb.RunInformation{ - RunNumber: 321500, + ut.CCDBService.On("GetRunInformation", uint64(321338)).Return(&ccdb.RunInformation{ + RunNumber: 321338, SOR: uint64(now + 7000), EOR: uint64(now + 10000), }, nil) @@ -642,9 +619,8 @@ func TestTrainingTaskService_UploadToCCDB_MissingExpectedFile(t *testing.T) { ut.TTRepo.AssertNotCalled(t, "Update", &tt) ut.TTRRepo.AssertCalled(t, "GetByType", ttId, models.Onnx) ut.FileService.AssertNotCalled(t, "OpenFile", "./local_file_temp.onnx") - ut.JAliEnService.AssertCalled(t, "ListAndParseDirectory", "/alice/sim/2024/LHC24f3/0") - ut.CCDBService.AssertCalled(t, "GetRunInformation", uint64(321000)) - ut.CCDBService.AssertCalled(t, "GetRunInformation", uint64(321500)) + ut.CCDBService.AssertCalled(t, "GetRunInformation", uint64(321321)) + ut.CCDBService.AssertCalled(t, "GetRunInformation", uint64(321338)) ut.CCDBService.AssertNotCalled(t, "UploadFile", uint64(now-10000), uint64(now+10000), "uploaded_file.onnx", file) } @@ -679,27 +655,17 @@ func TestTrainingTaskService_UploadToCCDB_ErrorReadingFile(t *testing.T) { ut.TTRepo.On("GetByID", ttId).Return(&tt, nil) ut.TTRepo.On("Update", &tt).Return(nil) ut.TTRRepo.On("GetByType", ttId, models.Onnx).Return(onnxFiles, nil) + ut.MonalisaService.On("GetRunList").Return(nil, errors.New("monalisa down")).Once() file := &mockReadCloser{} ut.FileService.On("OpenFile", "./local_file_temp.onnx").Return(nil, nil, errors.New("error reading file")) - ut.JAliEnService.On("ListAndParseDirectory", "/alice/sim/2024/LHC24f3/0").Return(&jalien.DirectoryContents{ - Subdirs: []jalien.Dir{ - {Name: "321000", Path: "/alice/sim/2024/LHC24f3/0/321000"}, - {Name: "321100", Path: "/alice/sim/2024/LHC24f3/0/321100"}, - {Name: "321321", Path: "/alice/sim/2024/LHC24f3/0/321321"}, - {Name: "321326", Path: "/alice/sim/2024/LHC24f3/0/321326"}, - {Name: "321338", Path: "/alice/sim/2024/LHC24f3/0/321338"}, - {Name: "321400", Path: "/alice/sim/2024/LHC24f3/0/321400"}, - {Name: "321500", Path: "/alice/sim/2024/LHC24f3/0/321500"}, - }, - }, nil) now := time.Now().UnixMilli() - ut.CCDBService.On("GetRunInformation", uint64(321000)).Return(&ccdb.RunInformation{ - RunNumber: 321000, + ut.CCDBService.On("GetRunInformation", uint64(321321)).Return(&ccdb.RunInformation{ + RunNumber: 321321, SOR: uint64(now - 10000), EOR: uint64(now - 9000), }, nil) - ut.CCDBService.On("GetRunInformation", uint64(321500)).Return(&ccdb.RunInformation{ - RunNumber: 321500, + ut.CCDBService.On("GetRunInformation", uint64(321338)).Return(&ccdb.RunInformation{ + RunNumber: 321338, SOR: uint64(now + 7000), EOR: uint64(now + 10000), }, nil) @@ -715,8 +681,7 @@ func TestTrainingTaskService_UploadToCCDB_ErrorReadingFile(t *testing.T) { ut.TTRepo.AssertNotCalled(t, "Update", &tt) ut.TTRRepo.AssertCalled(t, "GetByType", ttId, models.Onnx) ut.FileService.AssertCalled(t, "OpenFile", "./local_file_temp.onnx") - ut.JAliEnService.AssertCalled(t, "ListAndParseDirectory", "/alice/sim/2024/LHC24f3/0") - ut.CCDBService.AssertCalled(t, "GetRunInformation", uint64(321000)) - ut.CCDBService.AssertCalled(t, "GetRunInformation", uint64(321500)) + ut.CCDBService.AssertCalled(t, "GetRunInformation", uint64(321321)) + ut.CCDBService.AssertCalled(t, "GetRunInformation", uint64(321338)) ut.CCDBService.AssertNotCalled(t, "UploadFile", uint64(now-10000), uint64(now+10000), "uploaded_file.onnx", file) } diff --git a/web/templates/training-datasets/new.html b/web/templates/training-datasets/new.html index 9ce9806..9cdede7 100644 --- a/web/templates/training-datasets/new.html +++ b/web/templates/training-datasets/new.html @@ -92,61 +92,6 @@

Choose random subset

-
-
-
- - -
-
-
- - -
-
- - -
-
- - -
-
- -
-
- or - -
-
-