diff --git a/deploy/deploy.yaml b/deploy/deploy.yaml index d060a9fd..c800922c 100644 --- a/deploy/deploy.yaml +++ b/deploy/deploy.yaml @@ -545,7 +545,7 @@ lanes: env: HWLAB_WORKBENCH_RUNTIME_DB_URL: secretRef:hwlab-cloud-api-v03-db/database-url HWLAB_WORKBENCH_RUNTIME_PORT: "6671" - HWLAB_WORKBENCH_RUNTIME_DB_POOL_MAX: "8" + HWLAB_WORKBENCH_RUNTIME_DB_POOL_MAX: "32" HWLAB_WORKBENCH_RUNTIME_QUERY_TIMEOUT_MS: "10000" HWLAB_WORKBENCH_RUNTIME_READY_TIMEOUT_MS: "2000" OTEL_EXPORTER_OTLP_TRACES_ENDPOINT: http://otel-collector.platform-infra.svc.cluster.local:4318/v1/traces diff --git a/internal/workbenchruntime/service.go b/internal/workbenchruntime/service.go index f062cca9..0e5c3b1a 100644 --- a/internal/workbenchruntime/service.go +++ b/internal/workbenchruntime/service.go @@ -38,7 +38,10 @@ var workbenchRuntimeReadIndexes = []string{ "CREATE INDEX IF NOT EXISTS idx_workbench_sessions_updated ON workbench_sessions(updated_at DESC, session_id)", "CREATE INDEX IF NOT EXISTS idx_workbench_sessions_owner_updated ON workbench_sessions(owner_user_id, updated_at DESC, session_id)", "CREATE INDEX IF NOT EXISTS idx_workbench_sessions_status_updated ON workbench_sessions(status, updated_at DESC, session_id)", + "CREATE INDEX IF NOT EXISTS idx_workbench_messages_session_updated ON workbench_messages(session_id, updated_at ASC, message_id)", "CREATE INDEX IF NOT EXISTS idx_workbench_messages_session_role_updated ON workbench_messages(session_id, role, updated_at ASC, message_id)", + "CREATE INDEX IF NOT EXISTS idx_workbench_turns_trace_updated ON workbench_turns(trace_id, updated_at DESC)", + "CREATE INDEX IF NOT EXISTS idx_workbench_projection_checkpoints_trace_updated ON workbench_projection_checkpoints(trace_id, updated_at DESC)", } type factQuery struct { @@ -593,37 +596,70 @@ func (s *Server) queryMessageSummaries(ctx context.Context, sessionIDs []string) args = append(args, sid) placeholders = append(placeholders, fmt.Sprintf("$%d", len(args))) } - sqlText := fmt.Sprintf("SELECT session_id, COUNT(*) AS message_count, (ARRAY_AGG(message_json ORDER BY updated_at ASC) FILTER (WHERE role = 'user'))[1] AS first_user_message_json FROM workbench_messages WHERE session_id IN (%s) GROUP BY session_id", strings.Join(placeholders, ", ")) - result := []any{} - stage := "workbench_runtime.db.message_summaries" - err := s.withOtelInternalSpan(ctx, stage, s.dbQueryAttrs("workbench_messages", len(args), map[string]any{"db.index.expected": "idx_workbench_messages_session_updated|idx_workbench_messages_session_role_updated"}), func(spanCtx context.Context) error { - rows, err := s.db.QueryContext(spanCtx, sqlText, args...) + bySession := map[string]map[string]any{} + countSQL := fmt.Sprintf("SELECT session_id, COUNT(*) AS message_count FROM workbench_messages WHERE session_id IN (%s) GROUP BY session_id", strings.Join(placeholders, ", ")) + countStage := "workbench_runtime.db.message_counts" + err := s.withOtelInternalSpan(ctx, countStage, s.dbQueryAttrs("workbench_messages", len(args), map[string]any{"db.index.expected": "idx_workbench_messages_session_updated"}), func(spanCtx context.Context) error { + rows, err := s.db.QueryContext(spanCtx, countSQL, args...) if err != nil { return err } defer rows.Close() for rows.Next() { - var sessionID, firstUserMessageJSON sql.NullString + var sessionID sql.NullString var messageCount sql.NullInt64 - if err := rows.Scan(&sessionID, &messageCount, &firstUserMessageJSON); err != nil { + if err := rows.Scan(&sessionID, &messageCount); err != nil { return err } if !sessionID.Valid || strings.TrimSpace(sessionID.String) == "" { continue } item := map[string]any{"sessionId": sessionID.String, "messageCount": nullableInt(messageCount), "valuesRedacted": true} + bySession[sessionID.String] = item + } + return rows.Err() + }) + if err != nil { + return nil, wrapQueryStage(countStage, err) + } + firstUserSQL := fmt.Sprintf("SELECT DISTINCT ON (session_id) session_id, message_json FROM workbench_messages WHERE session_id IN (%s) AND role = 'user' ORDER BY session_id ASC, updated_at ASC, message_id ASC", strings.Join(placeholders, ", ")) + firstUserStage := "workbench_runtime.db.first_user_messages" + err = s.withOtelInternalSpan(ctx, firstUserStage, s.dbQueryAttrs("workbench_messages", len(args), map[string]any{"db.index.expected": "idx_workbench_messages_session_role_updated"}), func(spanCtx context.Context) error { + rows, err := s.db.QueryContext(spanCtx, firstUserSQL, args...) + if err != nil { + return err + } + defer rows.Close() + for rows.Next() { + var sessionID, firstUserMessageJSON sql.NullString + if err := rows.Scan(&sessionID, &firstUserMessageJSON); err != nil { + return err + } + if !sessionID.Valid || strings.TrimSpace(sessionID.String) == "" { + continue + } + item := bySession[sessionID.String] + if item == nil { + item = map[string]any{"sessionId": sessionID.String, "messageCount": 0, "valuesRedacted": true} + bySession[sessionID.String] = item + } if firstUserMessageJSON.Valid && strings.TrimSpace(firstUserMessageJSON.String) != "" { var message map[string]any if json.Unmarshal([]byte(firstUserMessageJSON.String), &message) == nil && len(message) > 0 { item["firstUserMessagePreview"] = firstUserPreview([]map[string]any{message}) } } - result = append(result, item) } return rows.Err() }) if err != nil { - return nil, wrapQueryStage(stage, err) + return nil, wrapQueryStage(firstUserStage, err) + } + result := []any{} + for _, sid := range cleaned { + if item := bySession[sid]; item != nil { + result = append(result, item) + } } return result, nil }