Back to home page

EIC code displayed by LXR

 
 

    


File indexing completed on 2026-10-01 09:17:03

0001 
0002 // Copyright 2020, Jefferson Science Associates, LLC.
0003 // Subject to the terms in the LICENSE file found in the top-level directory.
0004 
0005 #include "JBenchmarker.h"
0006 
0007 #include <JANA/Utils/JCpuInfo.h>
0008 
0009 #include <fstream>
0010 #include <cmath>
0011 #include <iomanip>
0012 #include <ios>
0013 #include <sys/stat.h>
0014 #include <iostream>
0015 #include <vector>
0016 
0017 JBenchmarker::JBenchmarker(JApplication* app) : m_app(app) {
0018 
0019     m_max_threads = JCpuInfo::GetNumCpus();
0020 
0021     auto params = app->GetJParameterManager();
0022 
0023     m_logger = params->GetLogger("benchmark");
0024 
0025     params->SetParameter("jana:nevents", 0);
0026     // Prevent users' choice of nevents from interfering with everything
0027 
0028     params->SetDefaultParameter(
0029             "benchmark:nsamples",
0030             m_nsamples,
0031             "Number of samples for each benchmark test");
0032 
0033     params->SetDefaultParameter(
0034             "benchmark:minthreads",
0035             m_min_threads,
0036             "Minimum number of threads for benchmark test");
0037 
0038     params->SetDefaultParameter(
0039             "benchmark:maxthreads",
0040             m_max_threads,
0041             "Maximum number of threads for benchmark test");
0042 
0043     params->SetDefaultParameter(
0044             "benchmark:use_log_scale",
0045             m_use_log_scale,
0046             "Use log scale (instead of linear)");
0047 
0048     if (m_use_log_scale) {
0049         // A thread step of 1 won't work for log scale, so in this case we default to 2
0050         m_thread_step = 2;
0051     }
0052 
0053     params->SetDefaultParameter(
0054             "benchmark:threadstep",
0055             m_thread_step,
0056             "Delta number of threads between each benchmark test");
0057 
0058     params->SetDefaultParameter(
0059             "benchmark:resultsdir",
0060             m_output_dir,
0061             "Output directory name for benchmark test results");
0062 
0063     params->SetDefaultParameter(
0064             "benchmark:rates_filename",
0065             m_rates_filename,
0066             "Filename for benchmark rates");
0067 
0068     params->SetDefaultParameter(
0069             "benchmark:samples_filename",
0070             m_samples_filename,
0071             "Filename for benchmark samples");
0072 
0073     params->SetDefaultParameter(
0074             "benchmark:copyscript",
0075             m_copy_script,
0076             "Copy plotting script to results dir");
0077 
0078 
0079     params->SetParameter("nthreads", m_max_threads);
0080     // Otherwise JApplication::Scale() doesn't scale up. This is an interesting bug. TODO: Remove me when fixed.
0081 }
0082 
0083 
0084 JBenchmarker::~JBenchmarker() {}
0085 
0086 
0087 void JBenchmarker::RunUntilFinished() {
0088 
0089     LOG_INFO(m_logger) << "Running benchmarker with the following settings:" << std::endl
0090                        << "    benchmark:minthreads = " << m_min_threads << std::endl
0091                        << "    benchmark:maxthreads = " << m_max_threads << std::endl
0092                        << "    benchmark:threadstep = " << m_thread_step << std::endl
0093                        << "    benchmark:use_log_scale = " << m_use_log_scale << std::endl
0094                        << "    benchmark:nsamples = " << m_nsamples << std::endl
0095                        << "    benchmark:resultsdir = " << m_output_dir << std::endl
0096                        << "    benchmark:rates_filename = " << m_rates_filename << std::endl
0097                        << "    benchmark:samples_filename = " << m_samples_filename << std::endl;
0098 
0099 
0100     mkdir(m_output_dir.c_str(), S_IRWXU | S_IRWXG | S_IROTH | S_IXOTH);
0101 
0102     std::ofstream samples_file(m_output_dir + "/" + m_samples_filename);
0103     samples_file << "# nthreads     rate" << std::endl;
0104 
0105     std::ofstream rates_file(m_output_dir + "/" + m_rates_filename);
0106     rates_file << "# nthreads  avg_rate       rms" << std::endl;
0107 
0108 
0109     std::vector<size_t> nthreads_space;
0110     if (m_use_log_scale) {
0111         for (size_t i=m_min_threads; i<m_max_threads; i *= m_thread_step) {
0112             nthreads_space.push_back(i);
0113         }
0114     }
0115     else {
0116         // Use linear scale
0117         for (size_t i=m_min_threads; i<m_max_threads; i += m_thread_step) {
0118             nthreads_space.push_back(i);
0119         }
0120     }
0121     if (nthreads_space.back() != m_max_threads) {
0122         nthreads_space.push_back(m_max_threads);
0123     }
0124 
0125     m_app->SetTicker(false);
0126     m_app->Run(false);
0127 
0128     // Wait for events to start flowing indicating the source is primed
0129     for (int i = 0; i < 5; i++) {
0130         LOG_INFO(m_logger) << "Waiting for event source to start producing ... rate: " << m_app->GetInstantaneousRate() << LOG_END;
0131         std::this_thread::sleep_for(std::chrono::milliseconds(1000));
0132         auto rate = m_app->GetInstantaneousRate();
0133         if (rate > 10.0) {
0134             LOG_INFO(m_logger) << "Rate: " << rate << "Hz   -  ready to begin test" << LOG_END;
0135             break;
0136         }
0137     }
0138 
0139     for (size_t nthreads: nthreads_space) {
0140         if (m_app->IsQuitting()) {
0141             break;
0142         }
0143 
0144         m_app->Scale(nthreads);
0145 
0146         // Loop for at most 60 seconds waiting for the number of threads to update
0147         for (int i = 0; i < 60; i++) {
0148             std::this_thread::sleep_for(std::chrono::milliseconds(1000));
0149             if (m_app->GetNThreads() == nthreads) break;
0150         }
0151 
0152         // Accumulate avg and rms rates for all samples for each nthreads
0153         double avg = 0;
0154         double rms = 0;
0155         double sum = 0;
0156         double sum2 = 0;
0157 
0158         for (uint32_t isample = 0; isample < m_nsamples && !m_app->IsQuitting(); isample++) {
0159             // Acquire mNsamples instantaneous rate measurements. The
0160             // GetInstantaneousRate method will only update every 0.5
0161             // seconds so we just wait for 1 second between samples to
0162             // ensure independent measurements.
0163             std::this_thread::sleep_for(std::chrono::milliseconds(1000));
0164             auto rate = m_app->GetInstantaneousRate();
0165 
0166             sum += rate;
0167             sum2 += rate * rate;
0168             double N = (double) (isample + 1);
0169             avg = sum / N; // Overwrite with updated value after each sample
0170             rms = sqrt((sum2 + N * avg * avg - 2.0 * avg * sum) / N);
0171 
0172             LOG_INFO(m_logger)
0173                 << std::setprecision(2) << std::fixed
0174                 << "nthreads=" << m_app->GetNThreads()
0175                 << "  rate=" << rate << "Hz"
0176                 << "  (avg = " << avg << " +/- " << rms / sqrt(N) << " Hz)" << LOG_END;
0177 
0178             // Write line in sample file
0179             samples_file << std::setw(7) << nthreads << " "
0180                          << std::setw(12) << std::setprecision(2) << std::fixed << rate << std::endl;
0181             samples_file.flush();
0182         }
0183 
0184         // Write line in rates file
0185         rates_file << std::setw(7) << nthreads << " "
0186                    << std::setw(12) << std::setprecision(2) << std::fixed << avg << " "
0187                    << std::setw(10) << std::setprecision(2) << std::fixed << rms << std::endl;
0188         rates_file.flush();
0189     }
0190 
0191     // Close files
0192     // Hopefully, because we called flush(), the files will be partially filled even if we are SIGKILLed
0193     // before we reach this point.
0194 
0195     samples_file.close();
0196     rates_file.close();
0197 
0198     if (m_copy_script) {
0199         copy_to_output_dir("${JANA_HOME}/bin/jana-plot-scaletest.py");
0200         LOG_INFO(m_logger)
0201             << "Testing finished. To view a plot of test results:\n"
0202             << "    cd " << m_output_dir
0203             << "\n    ./jana-plot-scaletest.py\n" << LOG_END;
0204     }
0205     else {
0206         LOG_INFO(m_logger) 
0207             << "Testing finished. To view a plot of test results:\n"
0208             << "    cd " << m_output_dir << "\n"
0209             << "    $JANA_HOME/bin/jana-plot-scaletest.py\n" << LOG_END;
0210     }
0211     m_app->Stop(true);
0212 }
0213 
0214 
0215 void JBenchmarker::copy_to_output_dir(std::string filename) {
0216 
0217     // Substitute environment variables in given filename
0218     std::string new_fname = filename;
0219     while (auto pos_start = new_fname.find("${") != new_fname.npos) {
0220         auto pos_end = new_fname.find("}", pos_start + 3);
0221         if (pos_end != new_fname.npos) {
0222 
0223             std::string envar_name = new_fname.substr(pos_start + 1, pos_end - pos_start - 1);
0224             LOG_DEBUG(m_logger) << "Looking for env var '" << envar_name
0225                                 << "'" << LOG_END;
0226 
0227             auto envar = getenv(envar_name.c_str());
0228             if (envar) {
0229                 new_fname.replace(pos_start - 1, pos_end + 2 - pos_start, envar);
0230             } else {
0231                 LOG_ERROR(m_logger) << "Environment variable '"
0232                                     << envar_name
0233                                     << "' not set. Cannot copy "
0234                                     << filename << LOG_END;
0235                 return;
0236             }
0237         } else {
0238             LOG_ERROR(m_logger) << "Error in string format: "
0239                                 << filename << LOG_END;
0240         }
0241     }
0242 
0243     // Extract filename without path
0244     std::string base_fname = new_fname;
0245     if (auto pos = base_fname.rfind("/")) base_fname.erase(0, pos);
0246     auto out_name = m_output_dir + "/" + base_fname;
0247 
0248     // Copy file
0249     LOG_INFO(m_logger) << "Copying " << new_fname << " -> " << m_output_dir << LOG_END;
0250     std::ifstream src(new_fname, std::ios::binary);
0251     std::ofstream dst(out_name, std::ios::binary);
0252     dst << src.rdbuf();
0253 
0254     // Change permissions to match source
0255     struct stat st;
0256     stat(new_fname.c_str(), &st);
0257     chmod(out_name.c_str(), st.st_mode);
0258 }
0259 
0260 
0261 
0262 
0263 
0264 
0265 
0266 
0267 
0268 
0269 
0270 
0271