package main import ( "context" "encoding/json" "fmt" "time" "github.com/google/uuid" "github.com/jackc/pgx/v5/pgxpool" ) const maxToolResultSize = 4096 type store struct { pool *pgxpool.Pool } func newStore(ctx context.Context, databaseURL string) (*store, error) { if databaseURL == "" { return nil, nil } pool, err := pgxpool.New(ctx, databaseURL) if err != nil { return nil, fmt.Errorf("connect db: %w", err) } if err := pool.Ping(ctx); err != nil { pool.Close() return nil, fmt.Errorf("ping db: %w", err) } return &store{pool: pool}, nil } func (s *store) close() { if s.pool != nil { s.pool.Close() } } type session struct { ID string `json:"id"` Title string `json:"title"` Actor string `json:"actor"` CreatedAt time.Time `json:"created_at"` LastActiveAt time.Time `json:"last_active_at"` } type message struct { ID string `json:"id"` SessionID string `json:"session_id"` Role string `json:"role"` Content json.RawMessage `json:"content"` CreatedAt time.Time `json:"created_at"` } func (s *store) createSession(ctx context.Context, title string) (*session, error) { if s == nil { return &session{ID: "ephemeral", Title: title, Actor: "agent:nomos"}, nil } var id string err := s.pool.QueryRow(ctx, `INSERT INTO agent_sessions (title, actor) VALUES ($1, 'agent:nomos') RETURNING id`, title).Scan(&id) if err != nil { return nil, err } return &session{ID: id, Title: title, Actor: "agent:nomos", CreatedAt: time.Now(), LastActiveAt: time.Now()}, nil } func (s *store) saveMessage(ctx context.Context, sessionID, role string, content json.RawMessage) error { if s == nil { return nil } _, err := s.pool.Exec(ctx, `INSERT INTO agent_messages (session_id, role, content) VALUES ($1, $2, $3)`, sessionID, role, truncateToolResults(content)) return err } func truncateToolResults(content json.RawMessage) json.RawMessage { var m map[string]any if err := json.Unmarshal(content, &m); err != nil { return content } toolCalls, ok := m["tool_calls"].([]any) if !ok || len(toolCalls) == 0 { return content } changed := false for i, raw := range toolCalls { tc, ok := raw.(map[string]any) if !ok { continue } if result, ok := tc["result"]; ok { resultJSON, _ := json.Marshal(result) if len(resultJSON) > maxToolResultSize { tc["result"] = string(resultJSON[:maxToolResultSize]) + fmt.Sprintf("...truncated (%d bytes total)", len(resultJSON)) toolCalls[i] = tc changed = true } } } if !changed { return content } m["tool_calls"] = toolCalls out, err := json.Marshal(m) if err != nil { return content } return out } func (s *store) touchSession(ctx context.Context, id string) { if s != nil { s.pool.Exec(ctx, `UPDATE agent_sessions SET last_active_at=now() WHERE id=$1`, id) } } func (s *store) listSessions(ctx context.Context) ([]session, error) { if s == nil { return nil, nil } rows, err := s.pool.Query(ctx, `SELECT id, title, actor, created_at, last_active_at FROM agent_sessions ORDER BY last_active_at DESC LIMIT 50`) if err != nil { return nil, err } defer rows.Close() var out []session for rows.Next() { var sess session if err := rows.Scan(&sess.ID, &sess.Title, &sess.Actor, &sess.CreatedAt, &sess.LastActiveAt); err != nil { return nil, err } out = append(out, sess) } return out, rows.Err() } func (s *store) getMessages(ctx context.Context, sessionID string) ([]message, error) { if s == nil { return nil, nil } rows, err := s.pool.Query(ctx, `SELECT id, session_id, role, content, created_at FROM agent_messages WHERE session_id=$1 ORDER BY created_at ASC`, sessionID) if err != nil { return nil, err } defer rows.Close() var out []message for rows.Next() { var m message if err := rows.Scan(&m.ID, &m.SessionID, &m.Role, &m.Content, &m.CreatedAt); err != nil { return nil, err } out = append(out, m) } return out, rows.Err() } func (s *store) deleteSession(ctx context.Context, id string) error { if s == nil { return nil } _, err := s.pool.Exec(ctx, `DELETE FROM agent_messages WHERE session_id = $1`, id) if err != nil { return err } _, err = s.pool.Exec(ctx, `DELETE FROM agent_sessions WHERE id = $1`, id) return err } func (s *store) updateSessionTitle(ctx context.Context, id, title string) error { if s == nil { return nil } _, err := s.pool.Exec(ctx, `UPDATE agent_sessions SET title = $1 WHERE id = $2`, title, id) return err } // resolveAgentID looks up the UUID of the agent entity (e.g. "agent:nomos"). // Returns uuid.Nil if the store is absent or the slug is unknown. func (s *store) resolveAgentID(ctx context.Context, slug string) uuid.UUID { if s == nil { return uuid.Nil } var id uuid.UUID if err := s.pool.QueryRow(ctx, `SELECT id FROM entities WHERE slug = $1`, slug).Scan(&id); err != nil { return uuid.Nil } return id } // logActivity records a tool call. agent_id is the agent entity UUID and is // NOT NULL in the schema, so we skip logging when it can't be resolved. // The (nullable) session_id column carries the conversation id. func (s *store) logActivity(ctx context.Context, agentID uuid.UUID, sessionID, toolName, inputSummary, outputSummary string, durationMs int, success bool, correlationID string) { if s == nil || agentID == uuid.Nil { return } s.pool.Exec(ctx, ` INSERT INTO agent_activity (agent_id, session_id, activity_type, tool_name, input_summary, output_summary, duration_ms, success, correlation_id) VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9)`, agentID, sessionID, "tool_call", toolName, inputSummary, outputSummary, durationMs, success, correlationID) }