• Home
  • Features
  • Pricing
  • Docs
  • Announcements
  • Sign In

llnl / dftracer / 22198188930

19 Feb 2026 08:09PM UTC coverage: 41.315%. First build
22198188930

Pull #338

github

web-flow
Merge ffd26ba39 into 1d847837f
Pull Request #338: Fix aggregator for paper

4009 of 13347 branches covered (30.04%)

Branch coverage included in aggregate %.

256 of 293 new or added lines in 10 files covered. (87.37%)

3326 of 4407 relevant lines covered (75.47%)

1002.76 hits per line

Source File
Press 'n' to go to next uncovered line, 'b' for previous

51.03
/src/dftracer/core/buffer/buffer.cpp
1
#include <dftracer/core/buffer/buffer.h>
2
template <>
3
std::shared_ptr<dftracer::BufferManager>
4
    dftracer::Singleton<dftracer::BufferManager>::instance = nullptr;
5
template <>
6
bool dftracer::Singleton<dftracer::BufferManager>::stop_creating_instances =
7
    false;
8
namespace dftracer {
9

10
void BufferManager::compress_and_write_if_needed(size_t size, bool force) {
1,065✔
11
  if (force || buffer_pos + size > this->config->write_buffer_size) {
1,065!
12
    // On forced flush, serialize any remaining aggregated data
13
    if (force && this->config->aggregation_enable) {
34!
NEW
14
      auto data = dftracer::AggregatedDataType();
×
NEW
15
      bool has_data = this->aggregator->get_previous_aggregations(data, false);
×
NEW
16
      if (has_data) {
×
NEW
17
        size_t agg_size = this->serializer->aggregated(
×
NEW
18
            buffer + buffer_pos + size, 0, rank, data);
×
NEW
19
        size += agg_size;
×
20
      }
NEW
21
    }
×
22

23
    if (this->config->compression) {
34✔
24
      size = this->compressor->compress(buffer, buffer_pos + size);
26✔
25
    } else {
26
      size = buffer_pos + size;
8✔
27
    }
28
    if (size > 0) {
34!
29
      size = this->writer->write(buffer, size, true);
34✔
30
    }
31
    buffer_pos = 0;
34✔
32
  } else {
33
    buffer_pos += size;
1,031✔
34
  }
35
}
1,065✔
36
int BufferManager::initialize(const char* filename, HashType hostname_hash) {
38✔
37
  DFTRACER_LOG_DEBUG("BufferManager.initialize %s %d", filename, hostname_hash);
30✔
38
  this->config =
39
      dftracer::Singleton<dftracer::ConfigurationManager>::get_instance();
38✔
40
  if (buffer == nullptr) {
38!
41
    buffer = (char*)malloc(this->config->write_buffer_size + 16 * 1024);
38✔
42
  }
43
  buffer_pos = 0;
38✔
44
  if (!buffer) {
38!
45
    DFTRACER_LOG_ERROR("BufferManager.BufferManager Failed to allocate buffer",
×
46
                       "");
47
  }
48
  this->writer = dftracer::Singleton<dftracer::STDIOWriter>::get_instance();
38✔
49
  this->writer->initialize(filename);
38✔
50
  this->serializer = dftracer::Singleton<dftracer::JsonLines>::get_instance();
38✔
51
  this->aggregator = dftracer::Singleton<dftracer::Aggregator>::get_instance();
38✔
52
  if (this->config->compression) {
38✔
53
    this->compressor =
54
        dftracer::Singleton<dftracer::ZlibCompression>::get_instance();
30✔
55
    this->compressor->initialize(this->config->write_buffer_size);
30✔
56
  }
57
  size_t size = this->serializer->initialize(buffer, hostname_hash);
38✔
58
  compress_and_write_if_needed(size);
38✔
59
  return 0;
38✔
60
}
61

62
int BufferManager::finalize(int index, ProcessID process_id, bool end_sym) {
34✔
63
  std::unique_lock<std::shared_mutex> lock(mtx);
34✔
64
  if (buffer) {
34!
65
    size_t size = 0;
34✔
66
    if (this->config->aggregation_enable) {
34!
67
      auto data = dftracer::AggregatedDataType();
×
68
      this->aggregator->get_previous_aggregations(data, true);
×
69
      size = this->serializer->aggregated(buffer + buffer_pos, index,
×
70
                                          process_id, data);
71
      this->aggregator->finalize();
×
72
    }
×
73
    auto end_size =
74
        this->serializer->finalize(buffer + buffer_pos + size, end_sym);
34✔
75
    compress_and_write_if_needed(size + end_size, true);
34✔
76

77
    if (this->config->compression) this->compressor->finalize();
34✔
78
    this->writer->finalize(index);
34✔
79
    free(buffer);
34✔
80
    buffer = nullptr;
34✔
81
    buffer_pos = 0;
34✔
82
  }
83
  return 0;
34✔
84
}
34✔
85

86
void BufferManager::log_data_event(int index, ConstEventNameType event_name,
751✔
87
                                   ConstEventNameType category,
88
                                   TimeResolution start_time,
89
                                   TimeResolution duration,
90
                                   dftracer::Metadata* metadata,
91
                                   ProcessID process_id, ThreadID tid) {
92
  std::unique_lock<std::shared_mutex> lock(mtx);
751✔
93
  DFTRACER_LOG_DEBUG("BufferManager.log_data_event %d", index);
386✔
94
  size_t size = 0;
751✔
95
  bool enable_tracing = true;
751✔
96
  if (this->config->aggregation_enable && strcmp(category, "dftracer") != 0) {
751!
97
    enable_tracing = false;
×
98
    auto aggregated_key =
99
        AggregatedKey{category, event_name, start_time,     duration,
100
                      tid,      metadata,   get_app_name(), &rank};
×
101
    if (this->config->aggregation_type ==
×
102
        AggregationType::AGGREGATION_TYPE_SELECTIVE) {
103
      enable_tracing = !this->aggregator->should_aggregate(&aggregated_key);
×
104
    }
105
    if (!enable_tracing) {
×
106
      // Accumulate data; returns true when time interval changes
NEW
107
      bool interval_changed = this->aggregator->aggregate(aggregated_key);
×
108
      // Serialize aggregated data when moving to new time interval
NEW
109
      if (interval_changed) {
×
110
        auto data = dftracer::AggregatedDataType();
×
NEW
111
        this->aggregator->get_previous_aggregations(data, false);
×
NEW
112
        if (!data.empty()) {
×
NEW
113
          size = this->serializer->aggregated(buffer + buffer_pos, index,
×
114
                                              process_id, data);
115
        }
116
      }
×
117
    }
118
  }
×
119
  if (enable_tracing) {
751!
120
    size =
121
        this->serializer->data(buffer + buffer_pos, index, event_name, category,
751✔
122
                               start_time, duration, metadata, process_id, tid);
123
  }
124
  compress_and_write_if_needed(size);
751✔
125
}
751✔
126

127
void BufferManager::log_counter_event(int index, ConstEventNameType name,
×
128
                                      ConstEventNameType category,
129
                                      TimeResolution start_time,
130
                                      ProcessID process_id, ThreadID thread_id,
131
                                      dftracer::Metadata* metadata) {
132
  std::unique_lock<std::shared_mutex> lock(mtx);
×
133
  DFTRACER_LOG_DEBUG("BufferManager.log_counter_event %d", index);
×
134
  size_t size =
135
      this->serializer->counter(buffer + buffer_pos, index, name, category,
×
136
                                start_time, process_id, thread_id, metadata);
137
  compress_and_write_if_needed(size);
×
138
}
×
139

140
void BufferManager::log_metadata_event(ConstEventNameType name,
242✔
141
                                       ConstEventNameType value,
142
                                       ConstEventNameType ph,
143
                                       ProcessID process_id, ThreadID tid,
144
                                       bool is_string) {
145
  std::unique_lock<std::shared_mutex> lock(mtx);
242✔
146
  DFTRACER_LOG_DEBUG("BufferManager.log_metadata_event %s", value);
198✔
147
  size_t size = this->serializer->metadata(buffer + buffer_pos, name, value, ph,
242✔
148
                                           process_id, tid, is_string);
149
  compress_and_write_if_needed(size);
242✔
150
}
242✔
151
}  // namespace dftracer
STATUS · Troubleshooting · Open an Issue · Sales · Support · CAREERS · ENTERPRISE · START FREE TRIAL · SCHEDULE DEMO
ANNOUNCEMENTS · TWITTER · TOS & SLA · Supported CI Services · What's a CI service? · Automated Testing

© 2026 Coveralls, Inc