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

llnl / dftracer / 22324588877

23 Feb 2026 08:59PM UTC coverage: 41.877%. First build
22324588877

Pull #338

github

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

3924 of 12861 branches covered (30.51%)

Branch coverage included in aggregate %.

269 of 307 new or added lines in 11 files covered. (87.62%)

3294 of 4375 relevant lines covered (75.29%)

990.05 hits per line

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

52.4
/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,049✔
11
  if (force || buffer_pos + size > this->config->write_buffer_size) {
1,049!
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
      DFTRACER_LOG_DEBUG(
18✔
26
          "BufferManager.compress_and_write_if_needed compressed size %zu "
27
          "bytes",
28
          size);
29
    } else {
30
      size = buffer_pos + size;
8✔
31
    }
32
    if (size > 0) {
34!
33
      size = this->writer->write(buffer, size, true);
34✔
34
      DFTRACER_LOG_DEBUG(
26✔
35
          "BufferManager.compress_and_write_if_needed wrote %zu bytes", size);
36
    }
37
    buffer_pos = 0;
34✔
38
  } else {
39
    buffer_pos += size;
1,015✔
40
    DFTRACER_LOG_DEBUG(
598✔
41
        "BufferManager.compress_and_write_if_needed buffer_pos %zu not writing",
42
        buffer_pos);
43
  }
44
}
1,049✔
45
int BufferManager::initialize(const char* filename, HashType hostname_hash) {
38✔
46
  DFTRACER_LOG_DEBUG("BufferManager.initialize %s %d", filename, hostname_hash);
30✔
47
  this->config =
48
      dftracer::Singleton<dftracer::ConfigurationManager>::get_instance();
38✔
49
  if (buffer == nullptr) {
38!
50
    buffer = (char*)malloc(this->config->write_buffer_size + 16 * 1024);
34✔
51
  }
52
  buffer_pos = 0;
38✔
53
  if (!buffer) {
38!
54
    DFTRACER_LOG_ERROR("BufferManager.BufferManager Failed to allocate buffer",
×
55
                       "");
56
  }
57
  this->writer = dftracer::Singleton<dftracer::STDIOWriter>::get_instance();
38✔
58
  this->writer->initialize(filename);
38✔
59
  this->serializer = dftracer::Singleton<dftracer::JsonLines>::get_instance();
38✔
60
  this->aggregator = dftracer::Singleton<dftracer::Aggregator>::get_instance();
38✔
61
  if (this->config->compression) {
38✔
62
    this->compressor =
63
        dftracer::Singleton<dftracer::ZlibCompression>::get_instance();
30✔
64
    this->compressor->initialize(this->config->write_buffer_size);
30✔
65
  }
66
  size_t size = this->serializer->initialize(buffer, hostname_hash);
38✔
67
  compress_and_write_if_needed(size);
38✔
68
  return 0;
38✔
69
}
70

71
int BufferManager::finalize(int index, ProcessID process_id, bool end_sym) {
34✔
72
  std::unique_lock<std::shared_mutex> lock(mtx);
34✔
73
  if (buffer) {
34!
74
    size_t size = 0;
34✔
75
    if (this->config->aggregation_enable) {
34!
76
      auto data = dftracer::AggregatedDataType();
×
77
      this->aggregator->get_previous_aggregations(data, true);
×
78
      size = this->serializer->aggregated(buffer + buffer_pos, index,
×
79
                                          process_id, data);
80
      this->aggregator->finalize();
×
81
    }
×
82
    auto end_size =
83
        this->serializer->finalize(buffer + buffer_pos + size, end_sym);
34✔
84
    compress_and_write_if_needed(size + end_size, true);
34✔
85

86
    if (this->config->compression) this->compressor->finalize();
34✔
87
    this->writer->finalize(index);
34✔
88
    free(buffer);
34✔
89
    buffer = nullptr;
34✔
90
    buffer_pos = 0;
34✔
91
  }
92
  return 0;
34✔
93
}
34✔
94

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

NEW
125
          DFTRACER_LOG_DEBUG(
×
126
              "BufferManager.log_data_event serialized aggregated size %zu "
127
              "bytes",
128
              size);
129
        }
130
      }
×
131
    }
132
  }
×
133
  if (enable_tracing) {
751!
134
    size =
135
        this->serializer->data(buffer + buffer_pos, index, event_name, category,
751✔
136
                               start_time, duration, metadata, process_id, tid);
137
    DFTRACER_LOG_DEBUG(
386✔
138
        "BufferManager.log_data_event serialized tracing size %zu bytes", size);
139
  }
140
  compress_and_write_if_needed(size);
751✔
141
}
751✔
142

143
void BufferManager::log_counter_event(int index, ConstEventNameType name,
×
144
                                      ConstEventNameType category,
145
                                      TimeResolution start_time,
146
                                      ProcessID process_id, ThreadID thread_id,
147
                                      dftracer::Metadata* metadata) {
148
  std::unique_lock<std::shared_mutex> lock(mtx);
×
149
  DFTRACER_LOG_DEBUG("BufferManager.log_counter_event %d", index);
×
150
  size_t size =
151
      this->serializer->counter(buffer + buffer_pos, index, name, category,
×
152
                                start_time, process_id, thread_id, metadata);
153
  compress_and_write_if_needed(size);
×
154
}
×
155

156
void BufferManager::log_metadata_event(ConstEventNameType name,
226✔
157
                                       ConstEventNameType value,
158
                                       ConstEventNameType ph,
159
                                       ProcessID process_id, ThreadID tid,
160
                                       bool is_string) {
161
  std::unique_lock<std::shared_mutex> lock(mtx);
226✔
162
  DFTRACER_LOG_DEBUG("BufferManager.log_metadata_event %s", value);
182✔
163
  size_t size = this->serializer->metadata(buffer + buffer_pos, name, value, ph,
226✔
164
                                           process_id, tid, is_string);
165
  compress_and_write_if_needed(size);
226✔
166
}
226✔
167
}  // 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