Class: OpenC3::LogWriter
Overview
Creates a log. Can automatically cycle the log based on an elapsed time period or when the log file reaches a predefined size.
Direct Known Subclasses
Constant Summary collapse
- CYCLE_TIME_INTERVAL =
The cycle time interval. Cycle times are only checked at this level of granularity.
10- CLEANUP_DELAY =
Delay in seconds before trimming Redis streams
60- @@mutex =
Mutex protecting class variables
Mutex.new
- @@instances =
Array of instances used to keep track of cycling logs
[]
- @@cycle_thread =
Thread used to cycle logs across all log writers
nil- @@cycle_sleeper =
Sleeper used to delay cycle thread
nil
Instance Attribute Summary collapse
-
#cleanup_offsets ⇒ Object
Redis offsets for each topic to cleanup.
-
#cleanup_times ⇒ Object
Time at which to cleanup.
-
#cycle_hour ⇒ Object
readonly
Returns the value of attribute cycle_hour.
-
#cycle_minute ⇒ Object
readonly
Returns the value of attribute cycle_minute.
-
#cycle_size ⇒ Object
readonly
Returns the value of attribute cycle_size.
-
#cycle_time ⇒ Object
readonly
Returns the value of attribute cycle_time.
-
#filename ⇒ String
readonly
The filename of the packet log.
-
#logging_enabled ⇒ true/false
readonly
Whether logging is enabled.
-
#mutex ⇒ Mutex
readonly
Instance mutex protecting file.
-
#start_time ⇒ Time
readonly
Time that the current log file started.
Instance Method Summary collapse
- #bucket_filename ⇒ Object
- #cleanup ⇒ Object
-
#close_file(take_mutex = true) ⇒ Object
Closing a log file isn't critical so we just log an error.
-
#create_unique_filename(ext = extension) ⇒ Object
implementation details.
- #cycle_thread_body ⇒ Object
- #extension ⇒ Object
- #first_timestamp ⇒ Object
- #graceful_kill ⇒ Object
-
#initialize(remote_log_directory, logging_enabled = true, cycle_time = nil, cycle_size = 1_000_000_000, cycle_hour = nil, cycle_minute = nil, enforce_time_order = true, cycle_thread: true) ⇒ LogWriter
constructor
A new instance of LogWriter.
- #last_timestamp ⇒ Object
-
#prepare_write(time_nsec_since_epoch, data_length, redis_topic = nil, redis_offset = nil, allow_new_file: true, process_out_of_order: true, stored: false) ⇒ Object
process_out_of_order ignores the timestamps for the current entry (used to ignore timestamps on metadata entries, vs actual packets) stored: stored packets roll a new file when GOOD times move backward (a new / overlapping replay run).
-
#shutdown ⇒ Object
Stop all logging, close the current log file, and kill the logging threads.
-
#start ⇒ Object
Starts a new log file by closing the existing log file.
-
#start_new_file ⇒ Object
Starting a new log file is a critical operation so the entire method is wrapped with a rescue and handled with handle_critical_exception Assumes mutex has already been taken.
-
#stop ⇒ Object
Stops all logging and closes the current log file.
Constructor Details
#initialize(remote_log_directory, logging_enabled = true, cycle_time = nil, cycle_size = 1_000_000_000, cycle_hour = nil, cycle_minute = nil, enforce_time_order = true, cycle_thread: true) ⇒ LogWriter
Returns a new instance of LogWriter.
91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 |
# File 'lib/openc3/logs/log_writer.rb', line 91 def initialize( remote_log_directory, logging_enabled = true, cycle_time = nil, cycle_size = 1_000_000_000, cycle_hour = nil, cycle_minute = nil, enforce_time_order = true, cycle_thread: true ) @remote_log_directory = remote_log_directory @logging_enabled = ConfigParser.handle_true_false(logging_enabled) @cycle_time = ConfigParser.handle_nil(cycle_time) if @cycle_time @cycle_time = Integer(@cycle_time) raise "cycle_time must be >= #{CYCLE_TIME_INTERVAL}" if @cycle_time < CYCLE_TIME_INTERVAL end @cycle_size = ConfigParser.handle_nil(cycle_size) @cycle_size = Integer(@cycle_size) if @cycle_size @cycle_hour = ConfigParser.handle_nil(cycle_hour) @cycle_hour = Integer(@cycle_hour) if @cycle_hour @cycle_minute = ConfigParser.handle_nil(cycle_minute) @cycle_minute = Integer(@cycle_minute) if @cycle_minute @enforce_time_order = ConfigParser.handle_true_false(enforce_time_order) @out_of_order = false @bad_time = false @mutex = Mutex.new @file = nil @file_size = 0 @filename = nil @start_time = Time.now.utc @first_time = nil @last_time = nil @cancel_threads = false @last_offsets = {} @cleanup_offsets = [] @cleanup_times = [] @previous_time_nsec_since_epoch = nil # Stored packets older than 20 years are treated as invalid (e.g. 1970 before # clock/GPS sync). They never drive file segmentation nor become the ordering # baseline, so a late-syncing clock doesn't thrash log files. @min_valid_time_nsec = (@start_time - (20 * 365 * 24 * 60 * 60)).to_nsec_from_epoch # Validity (valid/invalid time) of the packets in the currently open file. Used to # segregate invalid-time (1970) packets into their own file. nil = no packets yet. @current_file_valid = nil @tmp_dir = Dir.mktmpdir @wait_threads = [] # This is an optimization to avoid creating a new entry object # each time we create an entry which we do a LOT! @entry = String.new if cycle_thread # Always make sure there is a cycle thread - (because it does trimming) @@mutex.synchronize do @@instances << self unless @@cycle_thread @@cycle_thread = OpenC3.safe_thread("Log cycle") do cycle_thread_body() end end end end end |
Instance Attribute Details
#cleanup_offsets ⇒ Object
Redis offsets for each topic to cleanup
56 57 58 |
# File 'lib/openc3/logs/log_writer.rb', line 56 def cleanup_offsets @cleanup_offsets end |
#cleanup_times ⇒ Object
Time at which to cleanup
59 60 61 |
# File 'lib/openc3/logs/log_writer.rb', line 59 def cleanup_times @cleanup_times end |
#cycle_hour ⇒ Object (readonly)
Returns the value of attribute cycle_hour.
43 44 45 |
# File 'lib/openc3/logs/log_writer.rb', line 43 def cycle_hour @cycle_hour end |
#cycle_minute ⇒ Object (readonly)
Returns the value of attribute cycle_minute.
47 48 49 |
# File 'lib/openc3/logs/log_writer.rb', line 47 def cycle_minute @cycle_minute end |
#cycle_size ⇒ Object (readonly)
Returns the value of attribute cycle_size.
38 39 40 |
# File 'lib/openc3/logs/log_writer.rb', line 38 def cycle_size @cycle_size end |
#cycle_time ⇒ Object (readonly)
Returns the value of attribute cycle_time.
34 35 36 |
# File 'lib/openc3/logs/log_writer.rb', line 34 def cycle_time @cycle_time end |
#filename ⇒ String (readonly)
Returns The filename of the packet log.
27 28 29 |
# File 'lib/openc3/logs/log_writer.rb', line 27 def filename @filename end |
#logging_enabled ⇒ true/false (readonly)
Returns Whether logging is enabled.
30 31 32 |
# File 'lib/openc3/logs/log_writer.rb', line 30 def logging_enabled @logging_enabled end |
#mutex ⇒ Mutex (readonly)
Returns Instance mutex protecting file.
53 54 55 |
# File 'lib/openc3/logs/log_writer.rb', line 53 def mutex @mutex end |
#start_time ⇒ Time (readonly)
Returns Time that the current log file started.
50 51 52 |
# File 'lib/openc3/logs/log_writer.rb', line 50 def start_time @start_time end |
Instance Method Details
#bucket_filename ⇒ Object
404 405 406 |
# File 'lib/openc3/logs/log_writer.rb', line 404 def bucket_filename "#{}__#{}" + extension end |
#cleanup ⇒ Object
184 185 186 187 188 189 |
# File 'lib/openc3/logs/log_writer.rb', line 184 def cleanup if @tmp_dir FileUtils.remove_entry_secure(@tmp_dir, true) @tmp_dir = nil end end |
#close_file(take_mutex = true) ⇒ Object
Closing a log file isn't critical so we just log an error. NOTE: This also trims the Redis stream to keep a full file's worth of data in the stream. This is what prevents continuous stream growth. Returns thread that moves log to bucket
354 355 356 357 358 359 360 361 362 363 364 365 366 367 368 369 370 371 372 373 374 375 376 377 378 379 380 381 382 383 384 385 386 387 388 389 390 391 392 393 394 395 396 397 398 399 400 401 402 |
# File 'lib/openc3/logs/log_writer.rb', line 354 def close_file(take_mutex = true) @mutex.lock if take_mutex # Remove old wait_threads to_remove = [] @wait_threads.each do |thread| unless thread.alive? to_remove << thread end end to_remove.each do |thread| @wait_threads.delete(thread) end begin if @file begin @file.close unless @file.closed? Logger.debug "Log File Closed : #{@filename}" # Only try to moce the file if we've written data to it # This is indicated by the first and last timestamps being set if @first_time and @last_time date = [0..7] # YYYYMMDD bucket_key = File.join(@remote_log_directory, date, bucket_filename()) # Cleanup timestamps here so they are unset for the next file @first_time = nil @last_time = nil @wait_threads << BucketUtilities.move_log_file_to_bucket(@filename, bucket_key) # Now that the file is in storage, trim the Redis stream after a delay @cleanup_offsets << {} @last_offsets.each do |redis_topic, last_offset| @cleanup_offsets[-1][redis_topic] = last_offset end @cleanup_times << (Time.now + CLEANUP_DELAY) @last_offsets.clear end rescue Exception => e Logger.error "Error closing #{@filename} : #{e.formatted}" end @file = nil @file_size = 0 @filename = nil end ensure @mutex.unlock if take_mutex end return @wait_threads end |
#create_unique_filename(ext = extension) ⇒ Object
implementation details
197 198 199 200 201 202 203 204 205 206 207 208 209 210 211 212 |
# File 'lib/openc3/logs/log_writer.rb', line 197 def create_unique_filename(ext = extension) # Create a filename that doesn't exist attempt = nil while true filename_parts = [attempt] filename_parts.unshift @label if @label filename = File.join(@tmp_dir, File.([@label, attempt], ext)) if File.exist?(filename) attempt ||= 0 attempt += 1 Logger.warn("Unexpected file name conflict: #{filename}") else return filename end end end |
#cycle_thread_body ⇒ Object
214 215 216 217 218 219 220 221 222 223 224 225 226 227 228 229 230 231 232 233 234 235 236 237 238 239 240 241 242 243 244 245 246 247 248 249 250 251 252 253 254 255 256 257 258 259 260 261 262 263 264 265 266 267 268 269 270 271 272 273 |
# File 'lib/openc3/logs/log_writer.rb', line 214 def cycle_thread_body @@cycle_sleeper = Sleeper.new while true start_time = Time.now @@mutex.synchronize do @@instances.each do |instance| # The check against start_time needs to be mutex protected to prevent a packet coming in between the check # and closing the file instance.mutex.synchronize do utc_now = Time.now.utc # Logger.debug("start:#{@start_time.to_f} now:#{utc_now.to_f} cycle:#{@cycle_time} new:#{(utc_now - @start_time) > @cycle_time}") if instance.logging_enabled and instance.filename # Logging and file opened # Cycle based on total time logging if (instance.cycle_time and (utc_now - instance.start_time) > instance.cycle_time) Logger.debug("Log writer start new file due to cycle time") instance.close_file(false) # Cycle daily at a specific time elsif (instance.cycle_hour and instance.cycle_minute and utc_now.hour == instance.cycle_hour and utc_now.min == instance.cycle_minute and instance.start_time.yday != utc_now.yday) Logger.debug("Log writer start new file daily") instance.close_file(false) # Cycle hourly at a specific time elsif (instance.cycle_minute and not instance.cycle_hour and utc_now.min == instance.cycle_minute and instance.start_time.hour != utc_now.hour) Logger.debug("Log writer start new file hourly") instance.close_file(false) end end # Check for cleanup time indexes_to_clear = [] instance.cleanup_times.each_with_index do |cleanup_time, index| if cleanup_time <= utc_now # Now that the file is in S3, trim the Redis stream up until the previous file. # This keeps one minute of data in Redis instance.cleanup_offsets[index].each do |redis_topic, cleanup_offset| target_match = redis_topic.match(/__\{?([^}_]+)\}?__/) db_shard = target_match ? Store.db_shard_for_target(target_match[1]) : 0 Topic.trim_topic(redis_topic, cleanup_offset, db_shard: db_shard) end indexes_to_clear << index end end if indexes_to_clear.length > 0 indexes_to_clear.each do |index| instance.cleanup_offsets[index] = nil instance.cleanup_times[index] = nil end instance.cleanup_offsets.compact! instance.cleanup_times.compact! end end end end # Only check whether to cycle at a set interval run_time = Time.now - start_time sleep_time = CYCLE_TIME_INTERVAL - run_time sleep_time = 0 if sleep_time < 0 break if @@cycle_sleeper.sleep(sleep_time) end end |
#extension ⇒ Object
408 409 410 |
# File 'lib/openc3/logs/log_writer.rb', line 408 def extension '.log'.freeze end |
#first_timestamp ⇒ Object
412 413 414 |
# File 'lib/openc3/logs/log_writer.rb', line 412 def Time.from_nsec_from_epoch(@first_time). # "YYYYMMDDHHmmSSNNNNNNNNN" end |
#graceful_kill ⇒ Object
191 192 193 |
# File 'lib/openc3/logs/log_writer.rb', line 191 def graceful_kill @cancel_threads = true end |
#last_timestamp ⇒ Object
416 417 418 |
# File 'lib/openc3/logs/log_writer.rb', line 416 def Time.from_nsec_from_epoch(@last_time). # "YYYYMMDDHHmmSSNNNNNNNNN" end |
#prepare_write(time_nsec_since_epoch, data_length, redis_topic = nil, redis_offset = nil, allow_new_file: true, process_out_of_order: true, stored: false) ⇒ Object
process_out_of_order ignores the timestamps for the current entry (used to ignore timestamps on metadata entries, vs actual packets) stored: stored packets roll a new file when GOOD times move backward (a new / overlapping replay run). Packets with an invalid time (e.g. 1970 before clock/GPS sync) - realtime or stored - are segregated into their own file: the first valid packet after a run of 1970s rolls a fresh, properly-named file (and vice-versa).
305 306 307 308 309 310 311 312 313 314 315 316 317 318 319 320 321 322 323 324 325 326 327 328 329 330 331 332 333 334 335 336 337 338 339 340 341 342 343 344 345 346 347 348 349 |
# File 'lib/openc3/logs/log_writer.rb', line 305 def prepare_write(time_nsec_since_epoch, data_length, redis_topic = nil, redis_offset = nil, allow_new_file: true, process_out_of_order: true, stored: false) # An invalid time is a packet from before the clock/GPS synced (e.g. 1970). This # applies to both realtime and stored packets so neither taints a file's filename. packet_valid = !(process_out_of_order and time_nsec_since_epoch < @min_valid_time_nsec) # This check includes logging_enabled again because it might have changed since we acquired the mutex # Ensures new files based on size, and ensures always increasing time order in files if @logging_enabled if !@file Logger.debug("Log writer start new file because no file opened") start_new_file() if allow_new_file elsif @cycle_size and ((@file_size + data_length) > @cycle_size) Logger.debug("Log writer start new file due to cycle size #{@cycle_size}") start_new_file() if allow_new_file elsif process_out_of_order and !@current_file_valid.nil? and packet_valid != @current_file_valid and allow_new_file # Time validity changed (1970 <-> valid): segregate invalid-time packets into # their own file so valid files keep a clean, correctly-timed filename. Logger.debug("Log writer start new file due to time validity change") start_new_file() elsif process_out_of_order and @enforce_time_order and @previous_time_nsec_since_epoch and (@previous_time_nsec_since_epoch > time_nsec_since_epoch) if stored and packet_valid and allow_new_file # Good stored time went backward => new/overlapping replay run. Roll a new # file so each stored file stays monotonic (start_new_file resets @previous_time). Logger.debug("Log writer start new file due to out of order stored time") start_new_file() elsif !stored # Warning: Creating new files here can cause lots of files to be created if packets make it through out of order # Changed to just a warning to prevent file thrashing unless @out_of_order Logger.warn("Log writer out of order time detected (increase buffer depth?): #{Time.from_nsec_from_epoch(@previous_time_nsec_since_epoch)} #{Time.from_nsec_from_epoch(time_nsec_since_epoch)}") @out_of_order = true end end end end @last_offsets[redis_topic] = redis_offset if redis_topic and redis_offset # This is needed for the redis offset marker entry at the end of the log file if process_out_of_order @current_file_valid = packet_valid if packet_valid @previous_time_nsec_since_epoch = time_nsec_since_epoch elsif !@bad_time Logger.warn("Log writer segregating invalid packet time (before clock sync?): #{Time.from_nsec_from_epoch(time_nsec_since_epoch)}") @bad_time = true end end end |
#shutdown ⇒ Object
Stop all logging, close the current log file, and kill the logging threads.
171 172 173 174 175 176 177 178 179 180 181 182 |
# File 'lib/openc3/logs/log_writer.rb', line 171 def shutdown threads = stop() @@mutex.synchronize do @@instances.delete(self) if @@instances.length <= 0 @@cycle_sleeper.cancel if @@cycle_sleeper OpenC3.kill_thread(self, @@cycle_thread) if @@cycle_thread @@cycle_thread = nil end end return threads end |
#start ⇒ Object
Starts a new log file by closing the existing log file. New log files are not created until packets are written by #write so this does not immediately create a log file on the filesystem.
160 161 162 |
# File 'lib/openc3/logs/log_writer.rb', line 160 def start @mutex.synchronize { close_file(false); @logging_enabled = true } end |
#start_new_file ⇒ Object
Starting a new log file is a critical operation so the entire method is wrapped with a rescue and handled with handle_critical_exception Assumes mutex has already been taken
278 279 280 281 282 283 284 285 286 287 288 289 290 291 292 293 294 295 296 297 |
# File 'lib/openc3/logs/log_writer.rb', line 278 def start_new_file close_file(false) if @file # Start log file @filename = create_unique_filename() @file = File.new(@filename, 'wb') @file_size = 0 @start_time = Time.now.utc @out_of_order = false @first_time = nil @last_time = nil @previous_time_nsec_since_epoch = nil @current_file_valid = nil Logger.debug "Log File Opened : #{@filename}" rescue => e Logger.error "Error starting new log file: #{e.formatted}" @logging_enabled = false OpenC3.handle_critical_exception(e) end |
#stop ⇒ Object
Stops all logging and closes the current log file.
165 166 167 168 |
# File 'lib/openc3/logs/log_writer.rb', line 165 def stop @mutex.synchronize { close_file(false); @logging_enabled = false; } return @wait_threads end |