srdusr
aboutsummaryrefslogtreecommitdiffstats
path: root/internal/store/store.go
diff options
context:
space:
mode:
Diffstat (limited to 'internal/store/store.go')
-rw-r--r--internal/store/store.go61
1 files changed, 61 insertions, 0 deletions
diff --git a/internal/store/store.go b/internal/store/store.go
index 0e57c3c..e0f27d1 100644
--- a/internal/store/store.go
+++ b/internal/store/store.go
@@ -74,6 +74,17 @@ CREATE TABLE IF NOT EXISTS client_certs (
cert_pem BLOB NOT NULL,
key_pem BLOB NOT NULL
);
+
+CREATE TABLE IF NOT EXISTS ws_messages (
+ id INTEGER PRIMARY KEY AUTOINCREMENT,
+ entry_id INTEGER NOT NULL,
+ started_at INTEGER NOT NULL,
+ direction TEXT NOT NULL,
+ opcode INTEGER NOT NULL,
+ payload BLOB NOT NULL
+);
+
+CREATE INDEX IF NOT EXISTS ws_messages_entry_id ON ws_messages(entry_id);
`
// Store is a handle to the history database. Safe for concurrent use.
@@ -683,6 +694,56 @@ func (s *Store) DeleteClientCert(id int64) error {
return nil
}
+// WSMessage is one captured WebSocket frame, tagged to the history entry
+// of the upgrade request/response that started its connection - see
+// internal/proxy/websocket.go for why it's one row per frame rather than
+// per reassembled logical message.
+type WSMessage struct {
+ ID int64
+ EntryID int64
+ StartedAt time.Time
+ Direction string // "client_to_server" or "server_to_client"
+ Opcode int // RFC 6455 opcode: 1 text, 2 binary, 8 close, 9 ping, 10 pong
+ Payload []byte
+}
+
+// AddWSMessage stores one captured frame and returns its assigned ID.
+func (s *Store) AddWSMessage(m WSMessage) (int64, error) {
+ res, err := s.db.Exec(
+ `INSERT INTO ws_messages (entry_id, started_at, direction, opcode, payload) VALUES (?, ?, ?, ?, ?)`,
+ m.EntryID, m.StartedAt.UnixMilli(), m.Direction, m.Opcode, m.Payload,
+ )
+ if err != nil {
+ return 0, fmt.Errorf("add ws message: %w", err)
+ }
+ return res.LastInsertId()
+}
+
+// ListWSMessages returns every frame captured for entryID's WebSocket
+// connection, in the order they were sent.
+func (s *Store) ListWSMessages(entryID int64) ([]WSMessage, error) {
+ rows, err := s.db.Query(
+ `SELECT id, entry_id, started_at, direction, opcode, payload FROM ws_messages WHERE entry_id = ? ORDER BY id`,
+ entryID,
+ )
+ if err != nil {
+ return nil, fmt.Errorf("list ws messages: %w", err)
+ }
+ defer rows.Close()
+
+ var out []WSMessage
+ for rows.Next() {
+ var m WSMessage
+ var startedAt int64
+ if err := rows.Scan(&m.ID, &m.EntryID, &startedAt, &m.Direction, &m.Opcode, &m.Payload); err != nil {
+ return nil, fmt.Errorf("scan ws message row: %w", err)
+ }
+ m.StartedAt = time.UnixMilli(startedAt)
+ out = append(out, m)
+ }
+ return out, rows.Err()
+}
+
// DeleteEntry removes a single history entry and its search index row.
func (s *Store) DeleteEntry(id int64) error {
tx, err := s.db.Begin()