Class: SimpleAcp::Storage::PostgreSQL

Inherits:
Base
  • Object
show all
Defined in:
lib/simple_acp/storage/postgresql.rb

Overview

PostgreSQL-backed storage for persistent deployments.

Stores data in PostgreSQL using the Sequel gem. Automatically creates required tables on first use.

Examples:

storage = SimpleAcp::Storage::PostgreSQL.new(
  url: "postgres://localhost/simple_acp"
)
server = SimpleAcp::Server::Base.new(storage: storage)

Instance Method Summary collapse

Constructor Details

#initialize(options = {}) ⇒ PostgreSQL

Initialize PostgreSQL storage.

Parameters:

  • options (Hash) (defaults to: {}) —

    configuration options

Options Hash (options):

  • :db (Sequel::Database) —

    existing Sequel connection

  • :url (String) —

    database URL (default: $DATABASE_URL or localhost)

  • :skip_setup (Boolean) —

    skip automatic table creation

  • :host (String) —

    database host

  • :port (Integer) —

    database port

  • :database (String) —

    database name

  • :user (String) —

    database user

  • :password (String) —

    database password



27
28
29
30
31
# File 'lib/simple_acp/storage/postgresql.rb', line 27

def initialize(options = {})
  super
  @db = options[:db] || connect_db(options)
  setup_tables unless options[:skip_setup]
end

Instance Method Details

#add_event(run_id, event) ⇒ Object

See Also:



119
120
121
122
123
124
125
126
127
# File 'lib/simple_acp/storage/postgresql.rb', line 119

def add_event(run_id, event)
  @db[:acp_events].insert(
    run_id: run_id,
    event_type: event.type,
    data: event.to_json,
    created_at: Time.now
  )
  event
end

#clear! ⇒ void

This method returns an undefined value.

Clear all stored data.



156
157
158
159
160
# File 'lib/simple_acp/storage/postgresql.rb', line 156

def clear!
  @db[:acp_events].delete
  @db[:acp_runs].delete
  @db[:acp_sessions].delete
end

#close ⇒ Object

See Also:



142
143
144
# File 'lib/simple_acp/storage/postgresql.rb', line 142

def close
  @db.disconnect if @db.respond_to?(:disconnect)
end

#delete_run(run_id) ⇒ Object

See Also:



66
67
68
69
# File 'lib/simple_acp/storage/postgresql.rb', line 66

def delete_run(run_id)
  @db[:acp_events].where(run_id: run_id).delete
  @db[:acp_runs].where(run_id: run_id).delete
end

#delete_session(session_id) ⇒ Object



114
115
116
# File 'lib/simple_acp/storage/postgresql.rb', line 114

def delete_session(session_id)
  @db[:acp_sessions].where(id: session_id).delete
end

#get_events(run_id, limit: 100, offset: 0) ⇒ Object

See Also:



130
131
132
133
134
135
136
137
138
139
# File 'lib/simple_acp/storage/postgresql.rb', line 130

def get_events(run_id, limit: 100, offset: 0)
  rows = @db[:acp_events]
    .where(run_id: run_id)
    .order(:created_at)
    .limit(limit)
    .offset(offset)
    .all

  rows.map { |row| Models::Events.from_hash(JSON.parse(row[:data])) }
end

#get_run(run_id) ⇒ Object

See Also:



34
35
36
37
38
39
# File 'lib/simple_acp/storage/postgresql.rb', line 34

def get_run(run_id)
  row = @db[:acp_runs].where(run_id: run_id).first
  return nil unless row

  deserialize_run(row)
end

#get_session(session_id) ⇒ Object

See Also:



87
88
89
90
91
92
# File 'lib/simple_acp/storage/postgresql.rb', line 87

def get_session(session_id)
  row = @db[:acp_sessions].where(id: session_id).first
  return nil unless row

  deserialize_session(row)
end

#list_runs(agent_name: nil, session_id: nil, limit: 10, offset: 0) ⇒ Object

See Also:



72
73
74
75
76
77
78
79
80
81
82
83
84
# File 'lib/simple_acp/storage/postgresql.rb', line 72

def list_runs(agent_name: nil, session_id: nil, limit: 10, offset: 0)
  dataset = @db[:acp_runs]
  dataset = dataset.where(agent_name: agent_name) if agent_name
  dataset = dataset.where(session_id: session_id) if session_id

  total = dataset.count
  rows = dataset.order(Sequel.desc(:created_at)).limit(limit).offset(offset).all

  {
    runs: rows.map { |row| deserialize_run(row) },
    total: total
  }
end

#ping ⇒ Object

See Also:



147
148
149
150
151
# File 'lib/simple_acp/storage/postgresql.rb', line 147

def ping
  @db.test_connection
rescue StandardError
  false
end

#save_run(run) ⇒ Object

See Also:



42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
# File 'lib/simple_acp/storage/postgresql.rb', line 42

def save_run(run)
  data = {
    run_id: run.run_id,
    agent_name: run.agent_name,
    session_id: run.session_id,
    status: run.status,
    output: run.output.to_json,
    error: run.error&.to_json,
    await_request: run.await_request&.to_json,
    created_at: run.created_at,
    finished_at: run.finished_at,
    updated_at: Time.now
  }

  if @db[:acp_runs].where(run_id: run.run_id).count.positive?
    @db[:acp_runs].where(run_id: run.run_id).update(data)
  else
    @db[:acp_runs].insert(data)
  end

  run
end

#save_session(session) ⇒ Object

See Also:



95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
# File 'lib/simple_acp/storage/postgresql.rb', line 95

def save_session(session)
  data = {
    id: session.id,
    history: session.history.to_json,
    state: session.state&.to_json,
    updated_at: Time.now
  }

  if @db[:acp_sessions].where(id: session.id).count.positive?
    @db[:acp_sessions].where(id: session.id).update(data)
  else
    data[:created_at] = Time.now
    @db[:acp_sessions].insert(data)
  end

  session
end