Class: SmartBrain::EventStore::Postgres

Inherits:
Object
  • Object
show all
Defined in:
lib/smart_brain/event_store/postgres.rb

Overview

Postgres-backed EventStore. Drop-in replacement for EventStore::InMemory: writes the truth layer (sessions/turns/messages/refs) and reproduces the same return shapes so the runtime and diagnostics are backend-agnostic.

Instance Method Summary collapse

Constructor Details

#initialize(db:, clock:) ⇒ Postgres

Returns a new instance of Postgres.



12
13
14
15
# File 'lib/smart_brain/event_store/postgres.rb', line 12

def initialize(db:, clock:)
  @db = db
  @clock = clock
end

Instance Method Details

#all_turns(session_id:) ⇒ Object



94
95
96
97
98
99
100
101
102
103
104
105
# File 'lib/smart_brain/event_store/postgres.rb', line 94

def all_turns(session_id:)
  dataset = session_id.nil? ? db[:turns] : db[:turns].where(session_id: session_id)
  dataset.order(:session_id, :seq).map do |t|
    {
      id: t[:id],
      session_id: t[:session_id],
      seq: t[:seq],
      created_at: iso8601(t[:created_at]),
      turn_events: symbolize(t[:events_json])
    }
  end
end

#append_turn(session_id:, turn_events:, created_at:, domain_id: 'legacy') ⇒ Object



17
18
19
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
# File 'lib/smart_brain/event_store/postgres.rb', line 17

def append_turn(session_id:, turn_events:, created_at:, domain_id: 'legacy')
  db.transaction do
    ensure_session(session_id, domain_id)
    seq = next_seq(session_id)
    turn_id = SecureRandom.uuid
    messages = normalize_messages(turn_id: turn_id, messages: turn_events[:messages] || [], created_at: created_at)
    refs = normalize_refs(turn_id: turn_id, refs: turn_events[:refs] || [], created_at: created_at)
    normalized_events = turn_events.merge(messages: messages, refs: refs)

    db[:turns].insert(
      id: turn_id,
      session_id: session_id,
      seq: seq,
      created_at: time_from(created_at),
      events_json: Sequel.pg_jsonb(symbolizable(normalized_events))
    )
    insert_messages(turn_id, messages)
    insert_refs(turn_id, refs)

    {
      id: turn_id,
      session_id: session_id,
      seq: seq,
      created_at: iso8601(created_at),
      turn_events: normalized_events
    }
  end
end

#entity_frequencies(session_id:, window_turns:) ⇒ Object



77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
# File 'lib/smart_brain/event_store/postgres.rb', line 77

def entity_frequencies(session_id:, window_turns:)
  payloads = db[:turns].where(session_id: session_id)
                       .order(Sequel.desc(:seq))
                       .limit(window_turns)
                       .select(:events_json)
                       .all
  freq = Hash.new(0)
  payloads.each do |row|
    events = symbolize(row[:events_json])
    Array(events[:entities]).each do |entity|
      canonical = entity[:canonical] || entity[:name]
      freq[canonical.to_s.downcase] += 1
    end
  end
  freq
end

#recent_refs(session_id:, limit:) ⇒ Object



63
64
65
66
67
68
69
70
71
72
73
74
75
# File 'lib/smart_brain/event_store/postgres.rb', line 63

def recent_refs(session_id:, limit:)
  last_turn_ids = db[:turns].where(session_id: session_id).order(Sequel.desc(:seq)).limit(limit).select(:id)
  db[:refs].where(turn_id: last_turn_ids).order(:created_at).map do |r|
    {
      id: r[:id],
      turn_id: r[:turn_id],
      ref_type: r[:ref_type],
      ref_uri: r[:ref_uri],
      ref_meta_json: symbolize(r[:ref_meta_json]),
      created_at: iso8601(r[:created_at])
    }
  end
end

#recent_turns(session_id:, limit:) ⇒ Object



50
51
52
53
54
55
56
57
58
59
60
61
# File 'lib/smart_brain/event_store/postgres.rb', line 50

def recent_turns(session_id:, limit:)
  last_turn_ids = db[:turns].where(session_id: session_id).order(Sequel.desc(:seq)).limit(limit).select(:id)
  db[:messages].where(turn_id: last_turn_ids).order(:created_at).map do |m|
    {
      turn_id: m[:turn_id],
      message_id: m[:id],
      role: m[:role],
      content: m[:content],
      created_at: iso8601(m[:created_at])
    }
  end
end

#turns_count(session_id:) ⇒ Object



46
47
48
# File 'lib/smart_brain/event_store/postgres.rb', line 46

def turns_count(session_id:)
  db[:turns].where(session_id: session_id).count
end