Line data Source code
1 : /*! 2 : * \file esys/repo/progress/consoleobserver.cpp 3 : * \brief 4 : * 5 : * \cond 6 : * __legal_b__ 7 : * 8 : * Copyright (c) 2026 Michel Gillet 9 : * Distributed under the MIT License. 10 : * (See accompanying file LICENSE.txt or 11 : * copy at https://opensource.org/licenses/MIT) 12 : * 13 : * __legal_e__ 14 : * \endcond 15 : * 16 : */ 17 : 18 : #include "esys/repo/esysrepo_prec.h" 19 : #include "esys/repo/progress/consoleobserver.h" 20 : #include "esys/repo/progress/repoassignedevent.h" 21 : #include "esys/repo/progress/repodoneevent.h" 22 : #include "esys/repo/progress/repofailedevent.h" 23 : #include "esys/repo/progress/repoprogressevent.h" 24 : 25 : #include <iomanip> 26 : #include <sstream> 27 : #include <type_traits> 28 : #include <variant> 29 : 30 : namespace esys::repo::progress 31 : { 32 : 33 : namespace 34 : { 35 : 36 : constexpr int k_max_percentage = 100; 37 : 38 : } // namespace 39 : 40 36 : ConsoleObserver::ConsoleObserver() = default; 41 : 42 36 : ConsoleObserver::~ConsoleObserver() = default; 43 : 44 1747 : void ConsoleObserver::on_sync_event(const SyncEvent &event) 45 : { 46 1747 : std::visit( 47 1747 : [this](const auto &e) { 48 : using T = std::decay_t<decltype(e)>; 49 123 : if constexpr (std::is_same_v<T, RepoAssignedEvent>) apply_assigned(e); 50 1153 : else if constexpr (std::is_same_v<T, RepoProgressEvent>) apply_progress(e); 51 121 : else if constexpr (std::is_same_v<T, RepoDoneEvent>) apply_done(e); 52 1 : else if constexpr (std::is_same_v<T, RepoFailedEvent>) apply_failed(e); 53 : }, 54 : event); 55 1747 : } 56 : 57 36 : void ConsoleObserver::set_enabled(bool enabled) 58 : { 59 36 : std::lock_guard lock(m_mutex); 60 36 : m_enabled = enabled; 61 36 : } 62 : 63 0 : bool ConsoleObserver::get_enabled() const 64 : { 65 0 : std::lock_guard lock(m_mutex); 66 0 : return m_enabled; 67 0 : } 68 : 69 206 : std::string ConsoleObserver::format_status_line() const 70 : { 71 206 : std::lock_guard lock(m_mutex); 72 374 : if (!m_enabled || m_by_worker.empty()) return {}; 73 : 74 168 : std::ostringstream oss; 75 516 : for (const auto &entry : m_by_worker) append_repo_token(oss, entry.second); 76 168 : return oss.str(); 77 374 : } 78 : 79 0 : void ConsoleObserver::print_status_line(std::ostream &os) const 80 : { 81 0 : os << format_status_line(); 82 0 : } 83 : 84 123 : void ConsoleObserver::apply_assigned(const RepoAssignedEvent &event) 85 : { 86 123 : std::lock_guard lock(m_mutex); 87 123 : erase_repo(event.get_repo_index()); 88 123 : ActiveRepo repo; 89 123 : repo.m_repo_index = event.get_repo_index(); 90 123 : repo.m_phase = RepoPhase::PENDING; 91 123 : m_by_worker[event.get_worker_id()] = repo; 92 123 : } 93 : 94 1153 : void ConsoleObserver::apply_progress(const RepoProgressEvent &event) 95 : { 96 1153 : std::lock_guard lock(m_mutex); 97 1153 : ActiveRepo *repo = find_repo(event.get_repo_index()); 98 1153 : if (repo == nullptr) 99 : { 100 : // Progress before assign (or after clear): track under synthetic worker -1 101 0 : ActiveRepo created; 102 0 : created.m_repo_index = event.get_repo_index(); 103 0 : m_by_worker[-1 - static_cast<int>(event.get_repo_index())] = created; 104 0 : repo = find_repo(event.get_repo_index()); 105 0 : if (repo == nullptr) return; 106 0 : } 107 : 108 1153 : repo->m_phase = event.get_phase(); 109 1153 : repo->m_counters = event.get_counters(); 110 1153 : repo->m_percentage = percentage_from_counters(event.get_phase(), event.get_counters()); 111 1153 : repo->m_done = (event.get_phase() == RepoPhase::DONE); 112 1153 : } 113 : 114 121 : void ConsoleObserver::apply_done(const RepoDoneEvent &event) 115 : { 116 121 : std::lock_guard lock(m_mutex); 117 121 : erase_repo(event.get_repo_index()); 118 121 : } 119 : 120 1 : void ConsoleObserver::apply_failed(const RepoFailedEvent &event) 121 : { 122 1 : std::lock_guard lock(m_mutex); 123 1 : erase_repo(event.get_repo_index()); 124 1 : } 125 : 126 245 : void ConsoleObserver::erase_repo(std::size_t repo_index) 127 : { 128 245 : for (auto it = m_by_worker.begin(); it != m_by_worker.end();) 129 : { 130 482 : if (it->second.m_repo_index == repo_index) 131 122 : it = m_by_worker.erase(it); 132 : else 133 1087 : ++it; 134 : } 135 245 : } 136 : 137 1153 : ConsoleObserver::ActiveRepo *ConsoleObserver::find_repo(std::size_t repo_index) 138 : { 139 2160 : for (auto &entry : m_by_worker) 140 : { 141 2160 : if (entry.second.m_repo_index == repo_index) return &entry.second; 142 : } 143 : return nullptr; 144 : } 145 : 146 1153 : int ConsoleObserver::percentage_from_counters(RepoPhase phase, const TransferCounters &counters) 147 : { 148 1153 : if (phase == RepoPhase::RESOLVING && counters.get_total_deltas() > 0) 149 130 : return (k_max_percentage * counters.get_indexed_deltas()) / counters.get_total_deltas(); 150 : 151 1023 : if (counters.get_total_objects() > 0) 152 636 : return (k_max_percentage * counters.get_received_objects()) / counters.get_total_objects(); 153 : 154 387 : if (counters.get_total_checkout_steps() > 0) 155 0 : return (k_max_percentage * counters.get_checkout_steps()) / counters.get_total_checkout_steps(); 156 : 157 : return -1; 158 : } 159 : 160 348 : int ConsoleObserver::step_from_phase(RepoPhase phase) 161 : { 162 : // Maps high-level RepoPhase onto the classic "n/6" console slot. 163 : // Early sideband FetchSteps that collapse to RECEIVING all show as 5 (intentional). 164 348 : switch (phase) 165 : { 166 : case RepoPhase::NOT_SET: 167 : case RepoPhase::PENDING: return 0; 168 : case RepoPhase::RECEIVING: return 5; 169 : case RepoPhase::RESOLVING: 170 : case RepoPhase::CHECKOUT: return 6; 171 : case RepoPhase::DONE: 172 : case RepoPhase::FAILED: return 6; 173 : default: return 0; 174 : } 175 : } 176 : 177 348 : void ConsoleObserver::append_repo_token(std::ostream &os, const ActiveRepo &repo) 178 : { 179 348 : os << "[" << std::setw(2) << repo.m_repo_index << ": " << step_from_phase(repo.m_phase) << "/6 "; 180 348 : if (repo.m_done || (repo.m_percentage == k_max_percentage) || (repo.m_percentage < 0)) 181 319 : os << " "; 182 : else 183 29 : os << std::setw(2) << repo.m_percentage; 184 348 : os << "]"; 185 348 : } 186 : 187 : } // namespace esys::repo::progress