From 20f866e703000c2e710b77b940746b3a2fddb9ac Mon Sep 17 00:00:00 2001 From: Jin Hai Date: Sun, 19 Jul 2026 09:26:44 +0800 Subject: [PATCH] Go: refactor (#17072) Refactor stats --------- Signed-off-by: Jin Hai --- cmd/ragflow_server.go | 4 +- internal/engine/clickhouse/clickhouse.go | 96 ------------------- .../{stats_ee.go => clickhouse_ee.go} | 18 ++++ internal/server/server_ee.go | 2 +- 4 files changed, 21 insertions(+), 99 deletions(-) delete mode 100644 internal/engine/clickhouse/clickhouse.go rename internal/engine/clickhouse/{stats_ee.go => clickhouse_ee.go} (77%) diff --git a/cmd/ragflow_server.go b/cmd/ragflow_server.go index 101259364e..eff6320eb4 100644 --- a/cmd/ragflow_server.go +++ b/cmd/ragflow_server.go @@ -368,8 +368,8 @@ func main() { common.Warn("Failed to initialize server variables from Redis, using defaults", zap.String("error", err.Error())) } - ctx := context.Background() - if err = server.StartServer(ctx, serverName); err != nil { + ctx, cancel := context.WithCancel(context.Background()) + if err = server.StartServer(ctx, cancel, serverName); err != nil { common.Error("Failed to start EE server", err) os.Exit(1) } diff --git a/internal/engine/clickhouse/clickhouse.go b/internal/engine/clickhouse/clickhouse.go deleted file mode 100644 index 5fe5029980..0000000000 --- a/internal/engine/clickhouse/clickhouse.go +++ /dev/null @@ -1,96 +0,0 @@ -// -// Copyright 2026 The InfiniFlow Authors. All Rights Reserved. -// -// Licensed under the Apache License, Version 2.0 (the "License"); -// you may not use this file except in compliance with the License. -// You may obtain a copy of the License at -// -// http://www.apache.org/licenses/LICENSE-2.0 -// -// Unless required by applicable law or agreed to in writing, software -// distributed under the License is distributed on an "AS IS" BASIS, -// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. -// See the License for the specific language governing permissions and -// limitations under the License. -// - -package clickhouse - -import ( - "context" - "fmt" - "sync" - "time" - - "github.com/ClickHouse/clickhouse-go/v2" - "github.com/ClickHouse/clickhouse-go/v2/lib/driver" -) - -var ( - globalDriver *Driver - once sync.Once -) - -// GetDriver gets global ClickHouse client instance -func GetDriver() *Driver { - return globalDriver -} - -type Driver struct { - conn driver.Conn - config *clickhouse.Options -} - -func Init(ctx context.Context, host string, port int, user, password, database string) error { - var err error - once.Do(func() { - address := fmt.Sprintf("%s:%d", host, port) - - globalDriver = &Driver{} - globalDriver.config = &clickhouse.Options{ - Addr: []string{address}, - Auth: clickhouse.Auth{ - Database: database, - Username: user, - Password: password, - }, - DialTimeout: 5 * time.Second, - MaxOpenConns: 10, - MaxIdleConns: 5, - ConnMaxLifetime: time.Hour, - } - - globalDriver.conn, err = clickhouse.Open(globalDriver.config) - if err != nil { - return - } - - if err = globalDriver.conn.Ping(ctx); err != nil { - return - } - }) - return err -} - -func (d *Driver) Ping(ctx context.Context) error { - return d.conn.Ping(ctx) -} - -func (d *Driver) Status() (map[string]interface{}, error) { - stats := d.conn.Stats() - version, err := d.conn.ServerVersion() - if err != nil { - return nil, fmt.Errorf("failed to get server version: %w", err) - } - return map[string]interface{}{ - "max_open_conns": stats.MaxOpenConns, - "max_idle_conns": stats.MaxIdleConns, - "open": stats.Open, - "idle": stats.Idle, - "server_version": version.String(), - }, nil -} - -func Close() error { - return globalDriver.conn.Close() -} diff --git a/internal/engine/clickhouse/stats_ee.go b/internal/engine/clickhouse/clickhouse_ee.go similarity index 77% rename from internal/engine/clickhouse/stats_ee.go rename to internal/engine/clickhouse/clickhouse_ee.go index 76404f8cc2..7010af4d55 100644 --- a/internal/engine/clickhouse/stats_ee.go +++ b/internal/engine/clickhouse/clickhouse_ee.go @@ -18,11 +18,29 @@ package clickhouse import ( "ragflow/internal/common" + "sync" ) +var ( + globalDriver *Driver + once sync.Once +) + +type Driver struct { +} + +// GetDriver gets global stats client instance +func GetDriver() *Driver { + return globalDriver +} + func (d *Driver) CollectModelUsage(modelUsage *common.ModelUsage) error { //if modelUsage != nil { // common.Info("CollectModelUsage", zap.Any("modelUsage", modelUsage.String())) //} return nil } + +func (d *Driver) Status() (map[string]interface{}, error) { + return nil, nil +} diff --git a/internal/server/server_ee.go b/internal/server/server_ee.go index ee2dc1df90..c358bae904 100644 --- a/internal/server/server_ee.go +++ b/internal/server/server_ee.go @@ -37,7 +37,7 @@ func newEEServer() *EEServer { return &EEServer{} } -func StartServer(ctx context.Context, serverName string) error { +func StartServer(ctx context.Context, cancel context.CancelFunc, serverName string) error { if serverEE == nil { return errors.New("server EE is nil") }