Class: RubynCode::Teams::Mailbox

Inherits:
Object
  • Object
show all
Defined in:
lib/rubyn_code/teams/mailbox.rb

Overview

JSONL-based mailbox for inter-agent messaging backed by SQLite.

Messages are stored in the mailbox_messages table with structured JSON content. Each message tracks read/unread state per recipient. Supports structured data payloads and correlation IDs for request/response tracking.

Instance Method Summary collapse

Constructor Details

#initialize(db) ⇒ Mailbox



16
17
18
19
# File 'lib/rubyn_code/teams/mailbox.rb', line 16

def initialize(db)
  @db = db
  ensure_table!
end

Instance Method Details

#broadcast(from:, content:, all_names:) ⇒ Array<String>

Broadcasts a message from one agent to all other agents.



137
138
139
140
141
142
143
# File 'lib/rubyn_code/teams/mailbox.rb', line 137

def broadcast(from:, content:, all_names:)
  recipients = all_names.reject { |n| n == from }

  recipients.map do |recipient|
    send(from: from, to: recipient, content: content, message_type: 'broadcast')
  end
end

#find_by_correlation_id(correlation_id) ⇒ Array<Hash>

Finds all messages matching a correlation ID. Useful for tracking request/response chains.



118
119
120
121
122
123
124
125
126
127
128
129
# File 'lib/rubyn_code/teams/mailbox.rb', line 118

def find_by_correlation_id(correlation_id)
  rows = @db.query(
    "      SELECT id, payload, correlation_id, data FROM mailbox_messages\n      WHERE correlation_id = ?\n      ORDER BY created_at ASC\n    SQL\n    [correlation_id]\n  ).to_a\n\n  rows.map { |r| parse_message_row(r) }\nend\n",

#pending_for(name) ⇒ Array<Hash>

Returns unread messages for the given agent WITHOUT marking them as read. Used by IdlePoller to check for pending work without consuming messages.



150
151
152
153
154
155
156
157
158
159
160
161
# File 'lib/rubyn_code/teams/mailbox.rb', line 150

def pending_for(name)
  rows = @db.query(
    "      SELECT id, payload, correlation_id, data FROM mailbox_messages\n      WHERE recipient = ? AND read = 0\n      ORDER BY created_at ASC\n    SQL\n    [name]\n  ).to_a\n\n  rows.map { |r| parse_message_row(r) }\nend\n",

#read_inbox(name) ⇒ Array<Hash>

Reads all unread messages for the given agent and marks them as read.



88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
# File 'lib/rubyn_code/teams/mailbox.rb', line 88

def read_inbox(name)
  rows = @db.query(
    "      SELECT id, payload, correlation_id, data FROM mailbox_messages\n      WHERE recipient = ? AND read = 0\n      ORDER BY created_at ASC\n    SQL\n    [name]\n  ).to_a\n\n  return [] if rows.empty?\n\n  ids = rows.map { |r| r['id'] }\n  messages = rows.map { |r| parse_message_row(r) }\n\n  # Mark all fetched messages as read in a single statement\n  placeholders = ids.map { '?' }.join(', ')\n  @db.execute(\n    \"UPDATE mailbox_messages SET read = 1 WHERE id IN (\#{placeholders})\",\n    ids\n  )\n\n  messages\nend\n",

#send(from:, to:, content:, message_type: 'message', correlation_id: nil, data: nil) ⇒ String

Sends a message from one agent to another.

rubocop:disable Metrics/ParameterLists



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
# File 'lib/rubyn_code/teams/mailbox.rb', line 31

def send(from:, to:, content:, message_type: 'message', correlation_id: nil, data: nil)
  id = SecureRandom.uuid
  now = Time.now.utc.iso8601

  payload = JSON.generate({
                            id: id,
                            from: from,
                            to: to,
                            content: content,
                            message_type: message_type,
                            timestamp: now
                          })

  data_json = data ? JSON.generate(data) : nil

  @db.execute(
    "      INSERT INTO mailbox_messages (id, sender, recipient, message_type, payload, correlation_id, data, read, created_at)\n      VALUES (?, ?, ?, ?, ?, ?, ?, 0, ?)\n    SQL\n    [id, from, to, message_type, payload, correlation_id, data_json, now]\n  )\n\n  id\nend\n",

#send_structured(from:, to:, type:, data:, content: nil, correlation_id: nil) ⇒ String

Sends a structured message with typed data payload. Convenience wrapper around #send for machine-to-machine communication.

rubocop:disable Metrics/ParameterLists



69
70
71
72
73
74
75
76
77
78
79
80
81
# File 'lib/rubyn_code/teams/mailbox.rb', line 69

def send_structured(from:, to:, type:, data:, content: nil, correlation_id: nil)
  content ||= "#{type}: #{data.inspect}"[0, 200]
  correlation_id ||= SecureRandom.uuid

  send(
    from: from,
    to: to,
    content: content,
    message_type: type,
    correlation_id: correlation_id,
    data: data
  )
end

#unread_count(name) ⇒ Integer

Returns the count of unread messages for the given agent.



167
168
169
170
171
172
173
# File 'lib/rubyn_code/teams/mailbox.rb', line 167

def unread_count(name)
  rows = @db.query(
    'SELECT COUNT(*) AS cnt FROM mailbox_messages WHERE recipient = ? AND read = 0',
    [name]
  ).to_a
  rows.first&.fetch('cnt', 0) || 0
end