Class: Sentiero::Stores::File

Inherits:
Sentiero::Store show all
Defined in:
lib/sentiero/stores/file.rb

Constant Summary

Constants inherited from Sentiero::Store

Sentiero::Store::MAX_METADATA_KEYS, Sentiero::Store::MAX_METADATA_VALUE_SIZE, Sentiero::Store::PROBLEM_TITLE_MAX, Sentiero::Store::VALID_ID, Sentiero::Store::VALID_STATUS

Instance Attribute Summary

Attributes inherited from Sentiero::Store

#limits

Instance Method Summary collapse

Methods inherited from Sentiero::Store

#supports_event_aggregates?

Methods included from Sentiero::Store::SessionStore

#each_session_events, #event_counts_by_day, #event_type_counts

Constructor Details

#initialize(path:, limits: nil) ⇒ File

File-based store, single-process dev/test only. Each session is a directory: meta.json (timestamps, metadata) + window_id.jsonl (one JSON event per line).



14
15
16
17
18
# File 'lib/sentiero/stores/file.rb', line 14

def initialize(path:, limits: nil)
  @limits = limits
  @root = ::File.expand_path(path)
  FileUtils.mkdir_p(@root)
end

Instance Method Details

#clear!Object



294
295
296
297
# File 'lib/sentiero/stores/file.rb', line 294

def clear!
  FileUtils.rm_rf(Dir.glob(::File.join(@root, "*")))
  nil
end

#count_occurrences(problem_id, after: nil) ⇒ Object



223
224
225
226
227
228
229
230
# File 'lib/sentiero/stores/file.rb', line 223

def count_occurrences(problem_id, after: nil)
  validate_id!(problem_id)
  _problems, occurrences, _server_events = read_error_data
  list = occurrences[problem_id] || []
  return list.size unless after
  after_f = after.to_f
  list.count { |occ| occ["timestamp"].to_f > after_f }
end

#delete_session(session_id) ⇒ Object



148
149
150
151
152
153
154
155
156
157
158
159
# File 'lib/sentiero/stores/file.rb', line 148

def delete_session(session_id)
  validate_id!(session_id)
  dir = session_path(session_id)
  FileUtils.rm_rf(dir) if ::File.directory?(dir)

  update_error_data do |_problems, occurrences, server_events|
    occurrences.each_value { |list| list.reject! { |occ| occ["session_id"] == session_id } }
    server_events.reject! { |event| event["session_id"] == session_id }
  end

  nil
end

#delete_window(ref) ⇒ Object



161
162
163
164
165
166
167
168
169
170
171
172
173
# File 'lib/sentiero/stores/file.rb', line 161

def delete_window(ref)
  validate_window_ref!(ref)
  session_id = ref.session_id
  window_id = ref.window_id
  FileUtils.rm_f(window_path(session_id, window_id))

  # Only the session directory (meta.json + window files) is replay
  # data; error data lives in the shared root-level JSON files, so
  # this must not go through delete_session (which also erases those).
  remaining = list_window_ids(session_id)
  FileUtils.rm_rf(session_path(session_id)) if remaining.empty?
  nil
end

#get_events(ref, after: nil, limit: nil) ⇒ Object



119
120
121
122
123
124
125
126
127
128
129
130
131
# File 'lib/sentiero/stores/file.rb', line 119

def get_events(ref, after: nil, limit: nil)
  validate_window_ref!(ref)
  session_id = ref.session_id
  window_id = ref.window_id
  events = read_window_events(session_id, window_id).sort_by { |event| event["timestamp"].to_f }

  if after
    idx = events.index { |event| event["timestamp"].to_f > after.to_f }
    events = idx ? events[idx..] : []
  end

  limit ? events.first(limit) : events
end

#get_occurrences(problem_id, after: nil, limit: nil) ⇒ Object



214
215
216
217
218
219
220
221
# File 'lib/sentiero/stores/file.rb', line 214

def get_occurrences(problem_id, after: nil, limit: nil)
  validate_id!(problem_id)
  _problems, occurrences, _server_events = read_error_data
  list = occurrences[problem_id] || []
  result = list.sort_by { |occ| occ["timestamp"].to_f }
  result = result.select { |occ| occ["timestamp"].to_f > after.to_f } if after
  limit ? result.first(limit) : result
end

#get_problem(problem_id) ⇒ Object



208
209
210
211
212
# File 'lib/sentiero/stores/file.rb', line 208

def get_problem(problem_id)
  validate_id!(problem_id)
  problems, _occurrences, _server_events = read_error_data
  problems[problem_id]&.dup
end

#get_server_event(event_id) ⇒ Object



257
258
259
260
261
# File 'lib/sentiero/stores/file.rb', line 257

def get_server_event(event_id)
  validate_id!(event_id)
  _problems, _occurrences, server_events = read_error_data
  server_events.find { |e| e["id"] == event_id }
end

#get_session(session_id) ⇒ Object



97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
# File 'lib/sentiero/stores/file.rb', line 97

def get_session(session_id)
  validate_id!(session_id)
  meta = read_meta(session_id)
  return nil unless meta

  window_ids = list_window_ids(session_id)
  return nil if window_ids.empty?

  window_data = window_ids.map { |wid|
    events = read_window_events(session_id, wid)
    timestamps = events.filter_map { |event| event["timestamp"]&.to_f }
    window = {window_id: wid, event_count: events.size}
    if timestamps.any?
      window[:first_event_at] = timestamps.min
      window[:last_event_at] = timestamps.max
    end
    window
  }

  {session_id: session_id, windows: window_data}.merge(meta_fields(meta))
end

#list_problems(project:, limit:, offset: 0, status: nil, sort_by: nil, search: nil, since: nil, until_time: nil) ⇒ Object



193
194
195
196
197
198
199
200
201
202
203
204
205
206
# File 'lib/sentiero/stores/file.rb', line 193

def list_problems(project:, limit:, offset: 0, status: nil, sort_by: nil, search: nil, since: nil, until_time: nil)
  problems, _occurrences, _server_events = read_error_data
  filter_and_page_problems(
    problems.values,
    project: project,
    status: status,
    since: since,
    until_time: until_time,
    search: search,
    sort_by: sort_by,
    offset: offset,
    limit: limit
  )
end

#list_server_events(project:, limit:, name: nil, level: nil, session_id: nil, after: nil) ⇒ Object



263
264
265
266
267
268
269
270
271
272
273
274
# File 'lib/sentiero/stores/file.rb', line 263

def list_server_events(project:, limit:, name: nil, level: nil, session_id: nil, after: nil)
  _problems, _occurrences, server_events = read_error_data
  filter_server_events(
    server_events,
    project: project,
    name: name,
    level: level,
    session_id: session_id,
    after: after,
    limit: limit
  )
end

#list_sessions(limit:, offset: 0, since: nil, until_time: nil, sort_by: nil, search: nil, min_duration_ms: nil, min_events: nil) ⇒ Object



63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
# File 'lib/sentiero/stores/file.rb', line 63

def list_sessions(limit:, offset: 0, since: nil, until_time: nil, sort_by: nil, search: nil,
  min_duration_ms: nil, min_events: nil)
  session_ids = list_session_ids
  return [] if session_ids.empty?

  summaries = session_ids.filter_map { |sid| build_summary(sid) }

  if since
    since_f = since.to_f
    summaries = summaries.select { |summary| summary[:updated_at] >= since_f }
  end
  if until_time
    until_f = until_time.to_f
    summaries = summaries.select { |summary| summary[:updated_at] <= until_f }
  end
  if search && !search.empty?
    summaries = summaries.select { |summary| session_matches_search?(summary, search) }
  end
  if min_duration_ms || min_events
    summaries = summaries.select { |summary| session_meets_thresholds?(summary, min_duration_ms, min_events) }
  end

  case sort_by
  when "created_at"
    summaries.sort_by! { |summary| -summary[:created_at] }
  when "event_count"
    summaries.sort_by! { |summary| -summary[:event_count] }
  else
    summaries.sort_by! { |summary| -summary[:updated_at] }
  end

  summaries.slice(offset, limit) || []
end

#occurrences_for_session(session_id, limit: nil) ⇒ Object



276
277
278
279
280
# File 'lib/sentiero/stores/file.rb', line 276

def occurrences_for_session(session_id, limit: nil)
  validate_id!(session_id)
  _problems, occurrences, _server_events = read_error_data
  rows_for_session(occurrences.values.flatten, session_id, limit: limit)
end

#purge_older_than(seconds) ⇒ Object

Scan meta.json directly: the base list_sessions path is capped and newest-first, so it would never see the oldest (stale) sessions.



301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
# File 'lib/sentiero/stores/file.rb', line 301

def purge_older_than(seconds)
  cutoff = Time.now.to_f - seconds

  stale = list_session_ids.select { |sid|
    updated_at = read_meta(sid)&.fetch("updated_at", nil)
    updated_at && updated_at < cutoff
  }

  stale.each { |sid| delete_session(sid) }
  deleted = stale.size

  purge_error_data_older_than!(cutoff)

  deleted
end

#save_events(ref, events) ⇒ Object



20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
# File 'lib/sentiero/stores/file.rb', line 20

def save_events(ref, events)
  return if events.nil? || events.empty?

  validate_window_ref!(ref)
  session_id = ref.session_id
  window_id = ref.window_id

  now = Time.now.to_f
  session_dir = session_path(session_id)
  FileUtils.mkdir_p(session_dir)

  event_timestamps = events.filter_map { |event| event["timestamp"]&.to_f }
  batch_min = event_timestamps.min
  batch_max = event_timestamps.max

  ::File.open(window_path(session_id, window_id), "a") do |file|
    file.flock(::File::LOCK_EX)
    file.write(serialize_events(events))
  end

  update_meta(session_id) do |meta|
    if meta
      meta["updated_at"] = now
      meta["first_event_at"] = [meta["first_event_at"], batch_min].compact.min if batch_min
      meta["last_event_at"] = [meta["last_event_at"], batch_max].compact.max if batch_max
    else
      meta = {
        "created_at" => now,
        "updated_at" => now,
        "first_event_at" => batch_min,
        "last_event_at" => batch_max,
        "metadata" => nil
      }
    end
    meta
  end

  enforce_max_events(session_id)
  enforce_max_sessions

  nil
end

#save_metadata(session_id, metadata) ⇒ Object



133
134
135
136
137
138
139
140
141
142
143
144
145
146
# File 'lib/sentiero/stores/file.rb', line 133

def (session_id, )
  return unless .is_a?(Hash) && !.empty?
  validate_id!(session_id)
  validate_metadata!()

  update_meta(session_id) do |meta|
    next nil unless meta

    existing = meta["metadata"] || {}
    meta["metadata"] = existing.merge(.transform_keys(&:to_s))
    meta
  end
  nil
end

#save_occurrence(occurrence) ⇒ Object



175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
# File 'lib/sentiero/stores/file.rb', line 175

def save_occurrence(occurrence)
  validate_occurrence!(occurrence)
  fp = occurrence["fingerprint"]
  ts = occurrence["timestamp"].to_f
  occ_id = SecureRandom.uuid
  stored = occurrence.merge("id" => occ_id)

  update_error_data do |problems, occurrences, _server_events|
    existing = problems[fp]
    problems[fp] = existing ? touched_problem_attrs(existing, occurrence, ts) : new_problem_attrs(occurrence, ts)
    occurrences[fp] ||= []
    occurrences[fp] << stored
    evict_oldest_problems!(problems, occurrences, limits.max_problems)
  end
  (occurrence["session_id"], {"has_errors" => true}) if occurrence["session_id"]
  fp
end

#save_server_event(event) ⇒ Object



247
248
249
250
251
252
253
254
255
# File 'lib/sentiero/stores/file.rb', line 247

def save_server_event(event)
  validate_server_event!(event)
  stored = event.merge("id" => SecureRandom.uuid)
  update_error_data do |_problems, _occurrences, server_events|
    server_events << stored
    enforce_max_server_events!(server_events)
  end
  nil
end

#server_events_for_session(session_id, limit: nil) ⇒ Object



282
283
284
285
286
# File 'lib/sentiero/stores/file.rb', line 282

def server_events_for_session(session_id, limit: nil)
  validate_id!(session_id)
  _problems, _occurrences, server_events = read_error_data
  rows_for_session(server_events, session_id, limit: limit)
end

#session_ids_for_problem(problem_id, limit: nil) ⇒ Object



288
289
290
291
292
# File 'lib/sentiero/stores/file.rb', line 288

def session_ids_for_problem(problem_id, limit: nil)
  validate_id!(problem_id)
  _problems, occurrences, _server_events = read_error_data
  latest_session_ids(occurrences[problem_id] || [], limit: limit)
end

#update_problem_status(problem_id, status) ⇒ Object



232
233
234
235
236
237
238
239
240
241
242
243
244
245
# File 'lib/sentiero/stores/file.rb', line 232

def update_problem_status(problem_id, status)
  validate_id!(problem_id)
  validate_status!(status)
  update_error_data do |problems, _occurrences, _server_events|
    existing = problems[problem_id]
    next unless existing

    problems[problem_id] = existing.merge(
      status: status,
      resolved_at: (status == "resolved") ? Time.now.to_f : nil
    )
  end
  nil
end