fix(workbench-runtime): reduce sessions DB pool pressure (#1865)
This commit is contained in:
+1
-1
@@ -545,7 +545,7 @@ lanes:
|
|||||||
env:
|
env:
|
||||||
HWLAB_WORKBENCH_RUNTIME_DB_URL: secretRef:hwlab-cloud-api-v03-db/database-url
|
HWLAB_WORKBENCH_RUNTIME_DB_URL: secretRef:hwlab-cloud-api-v03-db/database-url
|
||||||
HWLAB_WORKBENCH_RUNTIME_PORT: "6671"
|
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_QUERY_TIMEOUT_MS: "10000"
|
||||||
HWLAB_WORKBENCH_RUNTIME_READY_TIMEOUT_MS: "2000"
|
HWLAB_WORKBENCH_RUNTIME_READY_TIMEOUT_MS: "2000"
|
||||||
OTEL_EXPORTER_OTLP_TRACES_ENDPOINT: http://otel-collector.platform-infra.svc.cluster.local:4318/v1/traces
|
OTEL_EXPORTER_OTLP_TRACES_ENDPOINT: http://otel-collector.platform-infra.svc.cluster.local:4318/v1/traces
|
||||||
|
|||||||
@@ -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_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_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_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_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 {
|
type factQuery struct {
|
||||||
@@ -593,37 +596,70 @@ func (s *Server) queryMessageSummaries(ctx context.Context, sessionIDs []string)
|
|||||||
args = append(args, sid)
|
args = append(args, sid)
|
||||||
placeholders = append(placeholders, fmt.Sprintf("$%d", len(args)))
|
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, ", "))
|
bySession := map[string]map[string]any{}
|
||||||
result := []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, ", "))
|
||||||
stage := "workbench_runtime.db.message_summaries"
|
countStage := "workbench_runtime.db.message_counts"
|
||||||
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 {
|
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, sqlText, args...)
|
rows, err := s.db.QueryContext(spanCtx, countSQL, args...)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
defer rows.Close()
|
defer rows.Close()
|
||||||
for rows.Next() {
|
for rows.Next() {
|
||||||
var sessionID, firstUserMessageJSON sql.NullString
|
var sessionID sql.NullString
|
||||||
var messageCount sql.NullInt64
|
var messageCount sql.NullInt64
|
||||||
if err := rows.Scan(&sessionID, &messageCount, &firstUserMessageJSON); err != nil {
|
if err := rows.Scan(&sessionID, &messageCount); err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
if !sessionID.Valid || strings.TrimSpace(sessionID.String) == "" {
|
if !sessionID.Valid || strings.TrimSpace(sessionID.String) == "" {
|
||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
item := map[string]any{"sessionId": sessionID.String, "messageCount": nullableInt(messageCount), "valuesRedacted": true}
|
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) != "" {
|
if firstUserMessageJSON.Valid && strings.TrimSpace(firstUserMessageJSON.String) != "" {
|
||||||
var message map[string]any
|
var message map[string]any
|
||||||
if json.Unmarshal([]byte(firstUserMessageJSON.String), &message) == nil && len(message) > 0 {
|
if json.Unmarshal([]byte(firstUserMessageJSON.String), &message) == nil && len(message) > 0 {
|
||||||
item["firstUserMessagePreview"] = firstUserPreview([]map[string]any{message})
|
item["firstUserMessagePreview"] = firstUserPreview([]map[string]any{message})
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
result = append(result, item)
|
|
||||||
}
|
}
|
||||||
return rows.Err()
|
return rows.Err()
|
||||||
})
|
})
|
||||||
if err != nil {
|
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
|
return result, nil
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user