diff options
Diffstat (limited to 'internal/store/store.go')
| -rw-r--r-- | internal/store/store.go | 61 |
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() |