Commit a68041f7 by Quaternijkon Committed by GitHub

feat(dashboard): add traffic flow sankey chart (#5465)

* feat(dashboard): add traffic flow sankey chart

Add dashboard flow APIs and a Sankey-based flow view with user, optional API key, model, and channel layers.\n\nReuse the dashboard VChart palette, add precise link/node tooltips and interactions, and cover filtering, layer ordering, color stability, and error states with tests.

* feat: build flow chart from quota data

---------

Co-authored-by: CaIon <i@caion.me>
parent f9e508bd
...@@ -10,6 +10,24 @@ import ( ...@@ -10,6 +10,24 @@ import (
"github.com/gin-gonic/gin" "github.com/gin-gonic/gin"
) )
func parseFlowQuotaTimeRange(c *gin.Context) (int64, int64, bool) {
startTimestamp, err := strconv.ParseInt(c.Query("start_timestamp"), 10, 64)
if err != nil || startTimestamp <= 0 {
common.ApiErrorMsg(c, "invalid start_timestamp")
return 0, 0, false
}
endTimestamp, err := strconv.ParseInt(c.Query("end_timestamp"), 10, 64)
if err != nil || endTimestamp <= 0 {
common.ApiErrorMsg(c, "invalid end_timestamp")
return 0, 0, false
}
if endTimestamp < startTimestamp {
common.ApiErrorMsg(c, "invalid time range")
return 0, 0, false
}
return startTimestamp, endTimestamp, true
}
func GetAllQuotaDates(c *gin.Context) { func GetAllQuotaDates(c *gin.Context) {
startTimestamp, _ := strconv.ParseInt(c.Query("start_timestamp"), 10, 64) startTimestamp, _ := strconv.ParseInt(c.Query("start_timestamp"), 10, 64)
endTimestamp, _ := strconv.ParseInt(c.Query("end_timestamp"), 10, 64) endTimestamp, _ := strconv.ParseInt(c.Query("end_timestamp"), 10, 64)
...@@ -66,3 +84,48 @@ func GetUserQuotaDates(c *gin.Context) { ...@@ -66,3 +84,48 @@ func GetUserQuotaDates(c *gin.Context) {
}) })
return return
} }
func GetAllFlowQuotaDates(c *gin.Context) {
startTimestamp, endTimestamp, ok := parseFlowQuotaTimeRange(c)
if !ok {
return
}
username := c.Query("username")
dates, err := model.GetFlowQuotaData(startTimestamp, endTimestamp, username, 0, c.GetInt("role"))
if err != nil {
common.ApiError(c, err)
return
}
c.JSON(http.StatusOK, gin.H{
"success": true,
"message": "",
"data": dates,
})
return
}
func GetUserFlowQuotaDates(c *gin.Context) {
userId := c.GetInt("id")
startTimestamp, endTimestamp, ok := parseFlowQuotaTimeRange(c)
if !ok {
return
}
if endTimestamp-startTimestamp > 2592000 {
c.JSON(http.StatusOK, gin.H{
"success": false,
"message": "时间跨度不能超过 1 个月",
})
return
}
dates, err := model.GetFlowQuotaData(startTimestamp, endTimestamp, "", userId, common.RoleCommonUser)
if err != nil {
common.ApiError(c, err)
return
}
c.JSON(http.StatusOK, gin.H{
"success": true,
"message": "",
"data": dates,
})
return
}
package controller
import (
"net/http"
"net/http/httptest"
"testing"
"github.com/QuantumNous/new-api/common"
"github.com/QuantumNous/new-api/model"
"github.com/gin-gonic/gin"
"github.com/stretchr/testify/require"
)
type flowQuotaResponse struct {
Success bool `json:"success"`
Message string `json:"message"`
Data []model.FlowQuotaData `json:"data"`
}
func setupFlowControllerTestDB(t *testing.T) {
t.Helper()
db := setupModelListControllerTestDB(t)
require.NoError(t, db.AutoMigrate(&model.Token{}, &model.QuotaData{}))
require.NoError(t, model.DB.Create(&model.Channel{Id: 1, Name: "east"}).Error)
require.NoError(t, model.DB.Create(&model.Token{Id: 11, UserId: 1, Key: "sk-primary", Name: "primary"}).Error)
require.NoError(t, model.DB.Create(&model.Token{Id: 22, UserId: 2, Key: "sk-backup", Name: "backup"}).Error)
require.NoError(t, model.DB.Create(&model.QuotaData{
UserID: 1,
Username: "alice",
NodeName: "node-a",
TokenID: 11,
UseGroup: "default",
ChannelID: 1,
ModelName: "gpt-a",
CreatedAt: 1100,
Count: 2,
Quota: 100,
TokenUsed: 40,
}).Error)
require.NoError(t, model.DB.Create(&model.QuotaData{
UserID: 2,
Username: "bob",
NodeName: "node-b",
TokenID: 22,
UseGroup: "vip",
ChannelID: 1,
ModelName: "gpt-b",
CreatedAt: 1200,
Count: 1,
Quota: 70,
TokenUsed: 30,
}).Error)
}
func decodeFlowQuotaResponse(t *testing.T, recorder *httptest.ResponseRecorder) flowQuotaResponse {
t.Helper()
require.Equal(t, http.StatusOK, recorder.Code)
var payload flowQuotaResponse
require.NoError(t, common.Unmarshal(recorder.Body.Bytes(), &payload))
require.True(t, payload.Success, payload.Message)
return payload
}
func TestGetAllFlowQuotaDatesUsesAdminDimensions(t *testing.T) {
setupFlowControllerTestDB(t)
recorder := httptest.NewRecorder()
ctx, _ := gin.CreateTestContext(recorder)
ctx.Set("role", common.RoleAdminUser)
ctx.Request = httptest.NewRequest(http.MethodGet, "/api/data/flow?start_timestamp=1000&end_timestamp=2000&username=bob", nil)
GetAllFlowQuotaDates(ctx)
payload := decodeFlowQuotaResponse(t, recorder)
require.Len(t, payload.Data, 1)
require.Equal(t, "bob", payload.Data[0].Username)
require.Equal(t, "vip", payload.Data[0].UseGroup)
require.Equal(t, "east", payload.Data[0].ChannelName)
require.Empty(t, payload.Data[0].TokenName)
require.Empty(t, payload.Data[0].NodeName)
}
func TestGetAllFlowQuotaDatesUsesRootDimensions(t *testing.T) {
setupFlowControllerTestDB(t)
recorder := httptest.NewRecorder()
ctx, _ := gin.CreateTestContext(recorder)
ctx.Set("role", common.RoleRootUser)
ctx.Request = httptest.NewRequest(http.MethodGet, "/api/data/flow?start_timestamp=1000&end_timestamp=2000&username=alice", nil)
GetAllFlowQuotaDates(ctx)
payload := decodeFlowQuotaResponse(t, recorder)
require.Len(t, payload.Data, 1)
require.Equal(t, "alice", payload.Data[0].Username)
require.Equal(t, "node-a", payload.Data[0].NodeName)
require.Equal(t, "primary", payload.Data[0].TokenName)
require.Equal(t, "default", payload.Data[0].UseGroup)
require.Equal(t, "east", payload.Data[0].ChannelName)
}
func TestGetUserFlowQuotaDatesRestrictsToAuthenticatedUser(t *testing.T) {
setupFlowControllerTestDB(t)
recorder := httptest.NewRecorder()
ctx, _ := gin.CreateTestContext(recorder)
ctx.Set("id", 1)
ctx.Request = httptest.NewRequest(http.MethodGet, "/api/data/flow/self?start_timestamp=1000&end_timestamp=2000", nil)
GetUserFlowQuotaDates(ctx)
payload := decodeFlowQuotaResponse(t, recorder)
require.Len(t, payload.Data, 1)
require.Empty(t, payload.Data[0].Username)
require.Equal(t, "primary", payload.Data[0].TokenName)
require.Equal(t, "default", payload.Data[0].UseGroup)
require.Empty(t, payload.Data[0].ChannelName)
}
func TestGetUserFlowQuotaDatesRejectsInvalidTimeRange(t *testing.T) {
setupFlowControllerTestDB(t)
recorder := httptest.NewRecorder()
ctx, _ := gin.CreateTestContext(recorder)
ctx.Set("id", 1)
ctx.Request = httptest.NewRequest(http.MethodGet, "/api/data/flow/self?start_timestamp=bad&end_timestamp=2000", nil)
GetUserFlowQuotaDates(ctx)
require.Equal(t, http.StatusOK, recorder.Code)
var payload flowQuotaResponse
require.NoError(t, common.Unmarshal(recorder.Body.Bytes(), &payload))
require.False(t, payload.Success)
require.Equal(t, "invalid start_timestamp", payload.Message)
}
...@@ -298,6 +298,7 @@ func RecordConsumeLog(c *gin.Context, userId int, params RecordConsumeLogParams) ...@@ -298,6 +298,7 @@ func RecordConsumeLog(c *gin.Context, userId int, params RecordConsumeLogParams)
username := c.GetString("username") username := c.GetString("username")
requestId := c.GetString(common.RequestIdKey) requestId := c.GetString(common.RequestIdKey)
upstreamRequestId := c.GetString(common.UpstreamRequestIdKey) upstreamRequestId := c.GetString(common.UpstreamRequestIdKey)
createdAt := common.GetTimestamp()
otherStr := common.MapToJsonStr(params.Other) otherStr := common.MapToJsonStr(params.Other)
// 判断是否需要记录 IP // 判断是否需要记录 IP
needRecordIp := false needRecordIp := false
...@@ -309,7 +310,7 @@ func RecordConsumeLog(c *gin.Context, userId int, params RecordConsumeLogParams) ...@@ -309,7 +310,7 @@ func RecordConsumeLog(c *gin.Context, userId int, params RecordConsumeLogParams)
log := &Log{ log := &Log{
UserId: userId, UserId: userId,
Username: username, Username: username,
CreatedAt: common.GetTimestamp(), CreatedAt: createdAt,
Type: LogTypeConsume, Type: LogTypeConsume,
Content: params.Content, Content: params.Content,
PromptTokens: params.PromptTokens, PromptTokens: params.PromptTokens,
...@@ -338,7 +339,18 @@ func RecordConsumeLog(c *gin.Context, userId int, params RecordConsumeLogParams) ...@@ -338,7 +339,18 @@ func RecordConsumeLog(c *gin.Context, userId int, params RecordConsumeLogParams)
} }
if common.DataExportEnabled { if common.DataExportEnabled {
gopool.Go(func() { gopool.Go(func() {
LogQuotaData(userId, username, params.ModelName, params.Quota, common.GetTimestamp(), params.PromptTokens+params.CompletionTokens) LogQuotaData(QuotaDataLogParams{
UserID: userId,
Username: username,
ModelName: params.ModelName,
Quota: params.Quota,
CreatedAt: createdAt,
TokenUsed: params.PromptTokens + params.CompletionTokens,
UseGroup: params.Group,
TokenID: params.TokenId,
ChannelID: params.ChannelId,
NodeName: common.NodeName,
})
}) })
} }
} }
...@@ -366,10 +378,11 @@ func RecordTaskBillingLog(params RecordTaskBillingLogParams) { ...@@ -366,10 +378,11 @@ func RecordTaskBillingLog(params RecordTaskBillingLogParams) {
tokenName = token.Name tokenName = token.Name
} }
} }
createdAt := common.GetTimestamp()
log := &Log{ log := &Log{
UserId: params.UserId, UserId: params.UserId,
Username: username, Username: username,
CreatedAt: common.GetTimestamp(), CreatedAt: createdAt,
Type: params.LogType, Type: params.LogType,
Content: params.Content, Content: params.Content,
TokenName: tokenName, TokenName: tokenName,
...@@ -384,6 +397,21 @@ func RecordTaskBillingLog(params RecordTaskBillingLogParams) { ...@@ -384,6 +397,21 @@ func RecordTaskBillingLog(params RecordTaskBillingLogParams) {
if err != nil { if err != nil {
common.SysLog("failed to record task billing log: " + err.Error()) common.SysLog("failed to record task billing log: " + err.Error())
} }
if params.LogType == LogTypeConsume && common.DataExportEnabled {
gopool.Go(func() {
LogQuotaData(QuotaDataLogParams{
UserID: params.UserId,
Username: username,
ModelName: params.ModelName,
Quota: params.Quota,
CreatedAt: createdAt,
UseGroup: params.Group,
TokenID: params.TokenId,
ChannelID: params.ChannelId,
NodeName: common.NodeName,
})
})
}
} }
func GetAllLogs(logType int, startTimestamp int64, endTimestamp int64, modelName string, username string, tokenName string, startIdx int, num int, channel int, group string, requestId string, upstreamRequestId string) (logs []*Log, total int64, err error) { func GetAllLogs(logType int, startTimestamp int64, endTimestamp int64, modelName string, username string, tokenName string, startIdx int, num int, channel int, group string, requestId string, upstreamRequestId string) (logs []*Log, total int64, err error) {
......
...@@ -40,6 +40,7 @@ func TestMain(m *testing.M) { ...@@ -40,6 +40,7 @@ func TestMain(m *testing.M) {
&Token{}, &Token{},
&Log{}, &Log{},
&Channel{}, &Channel{},
&QuotaData{},
&Ability{}, &Ability{},
&TopUp{}, &TopUp{},
&SubscriptionPlan{}, &SubscriptionPlan{},
...@@ -62,6 +63,7 @@ func truncateTables(t *testing.T) { ...@@ -62,6 +63,7 @@ func truncateTables(t *testing.T) {
DB.Exec("DELETE FROM tokens") DB.Exec("DELETE FROM tokens")
DB.Exec("DELETE FROM logs") DB.Exec("DELETE FROM logs")
DB.Exec("DELETE FROM channels") DB.Exec("DELETE FROM channels")
DB.Exec("DELETE FROM quota_data")
DB.Exec("DELETE FROM abilities") DB.Exec("DELETE FROM abilities")
DB.Exec("DELETE FROM top_ups") DB.Exec("DELETE FROM top_ups")
DB.Exec("DELETE FROM subscription_orders") DB.Exec("DELETE FROM subscription_orders")
......
...@@ -16,11 +16,28 @@ type QuotaData struct { ...@@ -16,11 +16,28 @@ type QuotaData struct {
Username string `json:"username" gorm:"index:idx_qdt_model_user_name,priority:2;size:64;default:''"` Username string `json:"username" gorm:"index:idx_qdt_model_user_name,priority:2;size:64;default:''"`
ModelName string `json:"model_name" gorm:"index:idx_qdt_model_user_name,priority:1;size:64;default:''"` ModelName string `json:"model_name" gorm:"index:idx_qdt_model_user_name,priority:1;size:64;default:''"`
CreatedAt int64 `json:"created_at" gorm:"bigint;index:idx_qdt_created_at,priority:2"` CreatedAt int64 `json:"created_at" gorm:"bigint;index:idx_qdt_created_at,priority:2"`
UseGroup string `json:"use_group" gorm:"index;size:64;default:''"`
TokenID int `json:"token_id" gorm:"index;default:0"`
ChannelID int `json:"channel_id" gorm:"index;default:0"`
NodeName string `json:"node_name" gorm:"index;size:64;default:''"`
TokenUsed int `json:"token_used" gorm:"default:0"` TokenUsed int `json:"token_used" gorm:"default:0"`
Count int `json:"count" gorm:"default:0"` Count int `json:"count" gorm:"default:0"`
Quota int `json:"quota" gorm:"default:0"` Quota int `json:"quota" gorm:"default:0"`
} }
type QuotaDataLogParams struct {
UserID int
Username string
ModelName string
Quota int
CreatedAt int64
TokenUsed int
UseGroup string
TokenID int
ChannelID int
NodeName string
}
func UpdateQuotaData() { func UpdateQuotaData() {
for { for {
if common.DataExportEnabled { if common.DataExportEnabled {
...@@ -34,34 +51,50 @@ func UpdateQuotaData() { ...@@ -34,34 +51,50 @@ func UpdateQuotaData() {
var CacheQuotaData = make(map[string]*QuotaData) var CacheQuotaData = make(map[string]*QuotaData)
var CacheQuotaDataLock = sync.Mutex{} var CacheQuotaDataLock = sync.Mutex{}
func logQuotaDataCache(userId int, username string, modelName string, quota int, createdAt int64, tokenUsed int) { func logQuotaDataCache(quotaData *QuotaData) {
key := fmt.Sprintf("%d-%s-%s-%d", userId, username, modelName, createdAt) key := fmt.Sprintf("%d\x00%s\x00%s\x00%d\x00%s\x00%d\x00%d\x00%s",
quotaData, ok := CacheQuotaData[key] quotaData.UserID,
quotaData.Username,
quotaData.ModelName,
quotaData.CreatedAt,
quotaData.UseGroup,
quotaData.TokenID,
quotaData.ChannelID,
quotaData.NodeName,
)
count := quotaData.Count
quota := quotaData.Quota
tokenUsed := quotaData.TokenUsed
cachedQuotaData, ok := CacheQuotaData[key]
if ok { if ok {
quotaData.Count += 1 cachedQuotaData.Count += count
quotaData.Quota += quota cachedQuotaData.Quota += quota
quotaData.TokenUsed += tokenUsed cachedQuotaData.TokenUsed += tokenUsed
} else { quotaData = cachedQuotaData
quotaData = &QuotaData{
UserID: userId,
Username: username,
ModelName: modelName,
CreatedAt: createdAt,
Count: 1,
Quota: quota,
TokenUsed: tokenUsed,
}
} }
CacheQuotaData[key] = quotaData CacheQuotaData[key] = quotaData
} }
func LogQuotaData(userId int, username string, modelName string, quota int, createdAt int64, tokenUsed int) { func LogQuotaData(params QuotaDataLogParams) {
// 只精确到小时 // 只精确到小时
createdAt = createdAt - (createdAt % 3600) createdAt := params.CreatedAt - (params.CreatedAt % 3600)
quotaData := &QuotaData{
UserID: params.UserID,
Username: params.Username,
ModelName: params.ModelName,
CreatedAt: createdAt,
UseGroup: params.UseGroup,
TokenID: params.TokenID,
ChannelID: params.ChannelID,
NodeName: params.NodeName,
Count: 1,
Quota: params.Quota,
TokenUsed: params.TokenUsed,
}
CacheQuotaDataLock.Lock() CacheQuotaDataLock.Lock()
defer CacheQuotaDataLock.Unlock() defer CacheQuotaDataLock.Unlock()
logQuotaDataCache(userId, username, modelName, quota, createdAt, tokenUsed) logQuotaDataCache(quotaData)
} }
func SaveQuotaDataCache() { func SaveQuotaDataCache() {
...@@ -74,13 +107,15 @@ func SaveQuotaDataCache() { ...@@ -74,13 +107,15 @@ func SaveQuotaDataCache() {
// 3. 如果没有数据,就插入数据 // 3. 如果没有数据,就插入数据
for _, quotaData := range CacheQuotaData { for _, quotaData := range CacheQuotaData {
quotaDataDB := &QuotaData{} quotaDataDB := &QuotaData{}
DB.Table("quota_data").Where("user_id = ? and username = ? and model_name = ? and created_at = ?", DB.Table("quota_data").
quotaData.UserID, quotaData.Username, quotaData.ModelName, quotaData.CreatedAt).First(quotaDataDB) Where("user_id = ? and username = ? and model_name = ? and created_at = ? and use_group = ? and token_id = ? and channel_id = ? and node_name = ?",
quotaData.UserID, quotaData.Username, quotaData.ModelName, quotaData.CreatedAt, quotaData.UseGroup, quotaData.TokenID, quotaData.ChannelID, quotaData.NodeName).
First(quotaDataDB)
if quotaDataDB.Id > 0 { if quotaDataDB.Id > 0 {
//quotaDataDB.Count += quotaData.Count //quotaDataDB.Count += quotaData.Count
//quotaDataDB.Quota += quotaData.Quota //quotaDataDB.Quota += quotaData.Quota
//DB.Table("quota_data").Save(quotaDataDB) //DB.Table("quota_data").Save(quotaDataDB)
increaseQuotaData(quotaData.UserID, quotaData.Username, quotaData.ModelName, quotaData.Count, quotaData.Quota, quotaData.CreatedAt, quotaData.TokenUsed) increaseQuotaData(quotaData)
} else { } else {
DB.Table("quota_data").Create(quotaData) DB.Table("quota_data").Create(quotaData)
} }
...@@ -89,12 +124,14 @@ func SaveQuotaDataCache() { ...@@ -89,12 +124,14 @@ func SaveQuotaDataCache() {
common.SysLog(fmt.Sprintf("保存数据看板数据成功,共保存%d条数据", size)) common.SysLog(fmt.Sprintf("保存数据看板数据成功,共保存%d条数据", size))
} }
func increaseQuotaData(userId int, username string, modelName string, count int, quota int, createdAt int64, tokenUsed int) { func increaseQuotaData(quotaData *QuotaData) {
err := DB.Table("quota_data").Where("user_id = ? and username = ? and model_name = ? and created_at = ?", err := DB.Table("quota_data").
userId, username, modelName, createdAt).Updates(map[string]interface{}{ Where("user_id = ? and username = ? and model_name = ? and created_at = ? and use_group = ? and token_id = ? and channel_id = ? and node_name = ?",
"count": gorm.Expr("count + ?", count), quotaData.UserID, quotaData.Username, quotaData.ModelName, quotaData.CreatedAt, quotaData.UseGroup, quotaData.TokenID, quotaData.ChannelID, quotaData.NodeName).
"quota": gorm.Expr("quota + ?", quota), Updates(map[string]interface{}{
"token_used": gorm.Expr("token_used + ?", tokenUsed), "count": gorm.Expr("count + ?", quotaData.Count),
"quota": gorm.Expr("quota + ?", quotaData.Quota),
"token_used": gorm.Expr("token_used + ?", quotaData.TokenUsed),
}).Error }).Error
if err != nil { if err != nil {
common.SysLog(fmt.Sprintf("increaseQuotaData error: %s", err)) common.SysLog(fmt.Sprintf("increaseQuotaData error: %s", err))
...@@ -104,14 +141,22 @@ func increaseQuotaData(userId int, username string, modelName string, count int, ...@@ -104,14 +141,22 @@ func increaseQuotaData(userId int, username string, modelName string, count int,
func GetQuotaDataByUsername(username string, startTime int64, endTime int64) (quotaData []*QuotaData, err error) { func GetQuotaDataByUsername(username string, startTime int64, endTime int64) (quotaData []*QuotaData, err error) {
var quotaDatas []*QuotaData var quotaDatas []*QuotaData
// 从quota_data表中查询数据 // 从quota_data表中查询数据
err = DB.Table("quota_data").Where("username = ? and created_at >= ? and created_at <= ?", username, startTime, endTime).Find(&quotaDatas).Error err = DB.Table("quota_data").
Select("user_id, username, model_name, created_at, sum(count) as count, sum(quota) as quota, sum(token_used) as token_used").
Where("username = ? and created_at >= ? and created_at <= ?", username, startTime, endTime).
Group("user_id, username, model_name, created_at").
Find(&quotaDatas).Error
return quotaDatas, err return quotaDatas, err
} }
func GetQuotaDataByUserId(userId int, startTime int64, endTime int64) (quotaData []*QuotaData, err error) { func GetQuotaDataByUserId(userId int, startTime int64, endTime int64) (quotaData []*QuotaData, err error) {
var quotaDatas []*QuotaData var quotaDatas []*QuotaData
// 从quota_data表中查询数据 // 从quota_data表中查询数据
err = DB.Table("quota_data").Where("user_id = ? and created_at >= ? and created_at <= ?", userId, startTime, endTime).Find(&quotaDatas).Error err = DB.Table("quota_data").
Select("user_id, username, model_name, created_at, sum(count) as count, sum(quota) as quota, sum(token_used) as token_used").
Where("user_id = ? and created_at >= ? and created_at <= ?", userId, startTime, endTime).
Group("user_id, username, model_name, created_at").
Find(&quotaDatas).Error
return quotaDatas, err return quotaDatas, err
} }
......
package model
import (
"fmt"
"github.com/QuantumNous/new-api/common"
"gorm.io/gorm"
)
type FlowQuotaData struct {
UserID int `json:"user_id,omitempty" gorm:"column:user_id"`
Username string `json:"username,omitempty" gorm:"column:username"`
NodeName string `json:"node_name,omitempty" gorm:"column:node_name"`
TokenID int `json:"token_id,omitempty" gorm:"column:token_id"`
TokenName string `json:"token_name,omitempty" gorm:"-"`
UseGroup string `json:"use_group" gorm:"column:use_group"`
ChannelID int `json:"channel_id,omitempty" gorm:"column:channel_id"`
ChannelName string `json:"channel_name,omitempty" gorm:"-"`
ModelName string `json:"model_name" gorm:"column:model_name"`
TokenUsed int `json:"token_used" gorm:"column:token_used"`
Count int `json:"count" gorm:"column:count"`
Quota int `json:"quota" gorm:"column:quota"`
}
func GetFlowQuotaData(startTime int64, endTime int64, username string, userID int, role int) ([]*FlowQuotaData, error) {
switch {
case role >= common.RoleRootUser:
return getRootFlowQuotaData(startTime, endTime, username)
case role >= common.RoleAdminUser:
return getAdminFlowQuotaData(startTime, endTime, username)
default:
return getSelfFlowQuotaData(startTime, endTime, userID)
}
}
func flowQuotaBaseQuery(startTime int64, endTime int64) *gorm.DB {
query := DB.Table("quota_data").
Where("use_group <> ''").
Where("created_at >= ? and created_at <= ?", startTime, endTime)
return query
}
func getSelfFlowQuotaData(startTime int64, endTime int64, userID int) ([]*FlowQuotaData, error) {
rows := make([]*FlowQuotaData, 0)
err := flowQuotaBaseQuery(startTime, endTime).
Select("token_id, use_group, model_name, sum(count) as count, sum(quota) as quota, sum(token_used) as token_used").
Where("user_id = ?", userID).
Group("token_id, use_group, model_name").
Order("quota DESC").
Find(&rows).Error
if err != nil {
return nil, err
}
return rows, fillFlowTokenNames(rows)
}
func getAdminFlowQuotaData(startTime int64, endTime int64, username string) ([]*FlowQuotaData, error) {
rows := make([]*FlowQuotaData, 0)
query := flowQuotaBaseQuery(startTime, endTime).
Select("user_id, username, use_group, model_name, channel_id, sum(count) as count, sum(quota) as quota, sum(token_used) as token_used")
if username != "" {
query = query.Where("username = ?", username)
}
err := query.
Group("user_id, username, use_group, model_name, channel_id").
Order("quota DESC").
Find(&rows).Error
if err != nil {
return nil, err
}
return rows, fillFlowChannelNames(rows)
}
func getRootFlowQuotaData(startTime int64, endTime int64, username string) ([]*FlowQuotaData, error) {
rows := make([]*FlowQuotaData, 0)
query := flowQuotaBaseQuery(startTime, endTime).
Select("user_id, username, node_name, token_id, use_group, model_name, channel_id, sum(count) as count, sum(quota) as quota, sum(token_used) as token_used")
if username != "" {
query = query.Where("username = ?", username)
}
err := query.
Group("user_id, username, node_name, token_id, use_group, model_name, channel_id").
Order("quota DESC").
Find(&rows).Error
if err != nil {
return nil, err
}
if err := fillFlowTokenNames(rows); err != nil {
return rows, err
}
return rows, fillFlowChannelNames(rows)
}
func fillFlowTokenNames(rows []*FlowQuotaData) error {
tokenIDSet := make(map[int]struct{})
tokenIDs := make([]int, 0)
for _, row := range rows {
if row.TokenID == 0 {
continue
}
if _, ok := tokenIDSet[row.TokenID]; ok {
continue
}
tokenIDSet[row.TokenID] = struct{}{}
tokenIDs = append(tokenIDs, row.TokenID)
}
if len(tokenIDs) == 0 {
return nil
}
var tokens []struct {
Id int `gorm:"column:id"`
Name string `gorm:"column:name"`
}
if err := DB.Model(&Token{}).Select("id, name").Where("id IN ?", tokenIDs).Find(&tokens).Error; err != nil {
return err
}
tokenNameByID := make(map[int]string, len(tokens))
for _, token := range tokens {
tokenNameByID[token.Id] = token.Name
}
// Deleted tokens are intentionally not resolved here: leave TokenName empty
// so the frontend can render a localized "deleted (id)" label instead.
for _, row := range rows {
if name := tokenNameByID[row.TokenID]; name != "" {
row.TokenName = name
}
}
return nil
}
func fillFlowChannelNames(rows []*FlowQuotaData) error {
channelIDSet := make(map[int]struct{})
channelIDs := make([]int, 0)
for _, row := range rows {
if row.ChannelID == 0 {
continue
}
if _, ok := channelIDSet[row.ChannelID]; ok {
continue
}
channelIDSet[row.ChannelID] = struct{}{}
channelIDs = append(channelIDs, row.ChannelID)
}
if len(channelIDs) == 0 {
return nil
}
channelNameByID := make(map[int]string, len(channelIDs))
if common.MemoryCacheEnabled {
for _, channelID := range channelIDs {
if channel, err := CacheGetChannel(channelID); err == nil {
channelNameByID[channelID] = channel.Name
}
}
} else {
var channels []struct {
Id int `gorm:"column:id"`
Name string `gorm:"column:name"`
}
if err := DB.Table("channels").Select("id, name").Where("id IN ?", channelIDs).Find(&channels).Error; err != nil {
return err
}
for _, channel := range channels {
channelNameByID[channel.Id] = channel.Name
}
}
for _, row := range rows {
if name := channelNameByID[row.ChannelID]; name != "" {
row.ChannelName = name
continue
}
if row.ChannelID > 0 {
row.ChannelName = fmt.Sprintf("channel-%d", row.ChannelID)
}
}
return nil
}
package model
import (
"testing"
"github.com/QuantumNous/new-api/common"
"github.com/stretchr/testify/require"
)
func seedFlowQuotaData(t *testing.T, quotaData QuotaData) {
t.Helper()
require.NoError(t, DB.Create(&quotaData).Error)
}
func seedFlowLookupData(t *testing.T) {
t.Helper()
require.NoError(t, DB.Create(&Channel{Id: 1, Name: "east"}).Error)
require.NoError(t, DB.Create(&Channel{Id: 2, Name: "west"}).Error)
require.NoError(t, DB.Create(&Token{Id: 11, UserId: 1, Key: "sk-primary", Name: "primary"}).Error)
require.NoError(t, DB.Create(&Token{Id: 22, UserId: 2, Key: "sk-backup", Name: "backup"}).Error)
require.NoError(t, DB.Delete(&Token{Id: 11}).Error)
}
func TestGetFlowQuotaDataUsesQuotaDataRoleSpecificDimensions(t *testing.T) {
truncateTables(t)
seedFlowLookupData(t)
seedFlowQuotaData(t, QuotaData{
UserID: 1,
Username: "alice",
NodeName: "node-a",
TokenID: 11,
UseGroup: "vip",
ModelName: "gpt-a",
ChannelID: 1,
CreatedAt: 1000,
Count: 2,
Quota: 100,
TokenUsed: 40,
})
seedFlowQuotaData(t, QuotaData{
UserID: 1,
Username: "alice",
NodeName: "node-a",
TokenID: 11,
UseGroup: "vip",
ModelName: "gpt-a",
ChannelID: 1,
CreatedAt: 1100,
Count: 1,
Quota: 50,
TokenUsed: 20,
})
seedFlowQuotaData(t, QuotaData{
UserID: 1,
Username: "alice",
NodeName: "node-a",
TokenID: 11,
UseGroup: "vip",
ModelName: "gpt-a",
ChannelID: 2,
CreatedAt: 1200,
Count: 1,
Quota: 25,
TokenUsed: 10,
})
seedFlowQuotaData(t, QuotaData{
UserID: 2,
Username: "bob",
NodeName: "node-b",
TokenID: 22,
UseGroup: "default",
ModelName: "gpt-b",
ChannelID: 1,
CreatedAt: 1300,
Count: 3,
Quota: 70,
TokenUsed: 30,
})
seedFlowQuotaData(t, QuotaData{
UserID: 1,
Username: "alice",
ModelName: "legacy",
CreatedAt: 1400,
Count: 99,
Quota: 999,
TokenUsed: 999,
})
rootRows, err := GetFlowQuotaData(900, 2000, "", 0, common.RoleRootUser)
require.NoError(t, err)
require.Len(t, rootRows, 3)
// Token 11 was soft-deleted, so its name is intentionally left empty for the
// frontend to render a localized "deleted (id)" label instead.
require.Equal(t, FlowQuotaData{
UserID: 1,
Username: "alice",
NodeName: "node-a",
TokenID: 11,
TokenName: "",
UseGroup: "vip",
ChannelID: 1,
ChannelName: "east",
ModelName: "gpt-a",
TokenUsed: 60,
Count: 3,
Quota: 150,
}, *rootRows[0])
// A token that still exists resolves to its current name.
require.Equal(t, 22, rootRows[1].TokenID)
require.Equal(t, "backup", rootRows[1].TokenName)
adminRows, err := GetFlowQuotaData(900, 2000, "alice", 0, common.RoleAdminUser)
require.NoError(t, err)
require.Len(t, adminRows, 2)
require.Equal(t, 0, adminRows[0].TokenID)
require.Empty(t, adminRows[0].TokenName)
require.Empty(t, adminRows[0].NodeName)
require.Equal(t, "alice", adminRows[0].Username)
require.Equal(t, "vip", adminRows[0].UseGroup)
require.Equal(t, "east", adminRows[0].ChannelName)
require.Equal(t, 150, adminRows[0].Quota)
selfRows, err := GetFlowQuotaData(900, 2000, "", 1, common.RoleCommonUser)
require.NoError(t, err)
require.Len(t, selfRows, 1)
require.Empty(t, selfRows[0].Username)
require.Equal(t, 0, selfRows[0].ChannelID)
require.Empty(t, selfRows[0].ChannelName)
require.Empty(t, selfRows[0].TokenName)
require.Equal(t, "vip", selfRows[0].UseGroup)
require.Equal(t, 175, selfRows[0].Quota)
}
func TestLogQuotaDataSplitsRowsByUseGroupTokenChannelAndNode(t *testing.T) {
truncateTables(t)
CacheQuotaDataLock.Lock()
CacheQuotaData = make(map[string]*QuotaData)
CacheQuotaDataLock.Unlock()
LogQuotaData(QuotaDataLogParams{
UserID: 1,
Username: "alice",
ModelName: "gpt-a",
CreatedAt: 3661,
UseGroup: "vip",
TokenID: 11,
ChannelID: 1,
NodeName: "node-a",
Quota: 100,
TokenUsed: 40,
})
LogQuotaData(QuotaDataLogParams{
UserID: 1,
Username: "alice",
ModelName: "gpt-a",
CreatedAt: 3700,
UseGroup: "vip",
TokenID: 11,
ChannelID: 1,
NodeName: "node-a",
Quota: 50,
TokenUsed: 20,
})
LogQuotaData(QuotaDataLogParams{
UserID: 1,
Username: "alice",
ModelName: "gpt-a",
CreatedAt: 3700,
UseGroup: "default",
TokenID: 11,
ChannelID: 1,
NodeName: "node-a",
Quota: 25,
TokenUsed: 10,
})
SaveQuotaDataCache()
var rows []QuotaData
require.NoError(t, DB.Order("quota DESC").Find(&rows).Error)
require.Len(t, rows, 2)
require.Equal(t, int64(3600), rows[0].CreatedAt)
require.Equal(t, "vip", rows[0].UseGroup)
require.Equal(t, 11, rows[0].TokenID)
require.Equal(t, 1, rows[0].ChannelID)
require.Equal(t, "node-a", rows[0].NodeName)
require.Equal(t, 2, rows[0].Count)
require.Equal(t, 150, rows[0].Quota)
require.Equal(t, 60, rows[0].TokenUsed)
require.Equal(t, "default", rows[1].UseGroup)
require.Equal(t, 25, rows[1].Quota)
}
...@@ -316,6 +316,8 @@ func SetApiRouter(router *gin.Engine) { ...@@ -316,6 +316,8 @@ func SetApiRouter(router *gin.Engine) {
dataRoute.GET("/", middleware.AdminAuth(), controller.GetAllQuotaDates) dataRoute.GET("/", middleware.AdminAuth(), controller.GetAllQuotaDates)
dataRoute.GET("/users", middleware.AdminAuth(), controller.GetQuotaDatesByUser) dataRoute.GET("/users", middleware.AdminAuth(), controller.GetQuotaDatesByUser)
dataRoute.GET("/self", middleware.UserAuth(), controller.GetUserQuotaDates) dataRoute.GET("/self", middleware.UserAuth(), controller.GetUserQuotaDates)
dataRoute.GET("/flow", middleware.AdminAuth(), controller.GetAllFlowQuotaDates)
dataRoute.GET("/flow/self", middleware.UserAuth(), controller.GetUserFlowQuotaDates)
logRoute.Use(middleware.CORS(), middleware.CriticalRateLimit()) logRoute.Use(middleware.CORS(), middleware.CriticalRateLimit())
{ {
......
...@@ -67,6 +67,11 @@ interface MultiSelectProps { ...@@ -67,6 +67,11 @@ interface MultiSelectProps {
*/ */
maxVisibleChips?: number maxVisibleChips?: number
/** /**
* Replaces individual chips with a compact summary while preserving the
* normal dropdown/search behaviour.
*/
renderSelectedSummary?: (values: string[]) => React.ReactNode
/**
* When true, clicking a chip's label copies its value to the clipboard * When true, clicking a chip's label copies its value to the clipboard
* instead of being inert. The remove (×) button keeps its own behaviour. * instead of being inert. The remove (×) button keeps its own behaviour.
*/ */
...@@ -258,6 +263,14 @@ export function MultiSelect(props: MultiSelectProps) { ...@@ -258,6 +263,14 @@ export function MultiSelect(props: MultiSelectProps) {
> >
<ComboboxValue> <ComboboxValue>
{(values: string[]) => { {(values: string[]) => {
if (props.renderSelectedSummary) {
return (
<span className='bg-muted text-muted-foreground flex h-[calc(--spacing(5.25))] w-fit items-center justify-center rounded-sm px-1.5 font-mono text-xs font-medium whitespace-nowrap'>
{props.renderSelectedSummary(values)}
</span>
)
}
const shouldLimit = const shouldLimit =
typeof props.maxVisibleChips === 'number' && !expanded typeof props.maxVisibleChips === 'number' && !expanded
const visibleValues = shouldLimit const visibleValues = shouldLimit
...@@ -327,7 +340,11 @@ export function MultiSelect(props: MultiSelectProps) { ...@@ -327,7 +340,11 @@ export function MultiSelect(props: MultiSelectProps) {
</ComboboxValue> </ComboboxValue>
<ComboboxChipsInput <ComboboxChipsInput
id={props.id} id={props.id}
placeholder={props.selected.length === 0 ? placeholder : undefined} placeholder={
props.selected.length === 0 && !props.renderSelectedSummary
? placeholder
: undefined
}
onKeyDown={handleKeyDown} onKeyDown={handleKeyDown}
aria-label={placeholder} aria-label={placeholder}
/> />
......
...@@ -17,7 +17,11 @@ along with this program. If not, see <https://www.gnu.org/licenses/>. ...@@ -17,7 +17,11 @@ along with this program. If not, see <https://www.gnu.org/licenses/>.
For commercial licensing, please contact support@quantumnous.com For commercial licensing, please contact support@quantumnous.com
*/ */
import { api } from '@/lib/api' import { api } from '@/lib/api'
import type { QuotaDataItem, UptimeGroupResult } from './types' import type {
FlowQuotaDataItem,
QuotaDataItem,
UptimeGroupResult,
} from './types'
// ============================================================================ // ============================================================================
// Dashboard APIs // Dashboard APIs
...@@ -61,6 +65,24 @@ export async function getUserQuotaDataByUsers(params: { ...@@ -61,6 +65,24 @@ export async function getUserQuotaDataByUsers(params: {
return res.data return res.data
} }
export async function getFlowQuotaDates(
params: {
start_timestamp: number
end_timestamp: number
default_time?: string
username?: string
},
isAdmin = false
) {
const endpoint = isAdmin ? '/api/data/flow' : '/api/data/flow/self'
const res = await api.get<{
success: boolean
data?: FlowQuotaDataItem[]
message?: string
}>(endpoint, { params })
return res.data
}
// Get uptime monitoring status for all services // Get uptime monitoring status for all services
export async function getUptimeStatus() { export async function getUptimeStatus() {
const res = await api.get<{ success: boolean; data: UptimeGroupResult[] }>( const res = await api.get<{ success: boolean; data: UptimeGroupResult[] }>(
......
...@@ -53,6 +53,8 @@ interface ModelsFilterProps { ...@@ -53,6 +53,8 @@ interface ModelsFilterProps {
preferences: DashboardChartPreferences preferences: DashboardChartPreferences
onFilterChange: (filters: DashboardFilters) => void onFilterChange: (filters: DashboardFilters) => void
onReset: () => void onReset: () => void
titleKey?: string
descriptionKey?: string
} }
/** /**
...@@ -145,8 +147,11 @@ export function ModelsFilter(props: ModelsFilterProps) { ...@@ -145,8 +147,11 @@ export function ModelsFilter(props: ModelsFilterProps) {
{t('Filter')} {t('Filter')}
</Button> </Button>
} }
title={t('Model Analytics Filters')} title={t(props.titleKey ?? 'Model Analytics Filters')}
description={t('Filter the model analytics view by time range and user.')} description={t(
props.descriptionKey ??
'Filter the model analytics view by time range and user.'
)}
contentClassName='max-sm:h-dvh max-sm:w-screen max-sm:max-w-none max-sm:rounded-none max-sm:p-4 sm:max-w-lg' contentClassName='max-sm:h-dvh max-sm:w-screen max-sm:max-w-none max-sm:rounded-none max-sm:p-4 sm:max-w-lg'
contentHeight='min(48vh, 460px)' contentHeight='min(48vh, 460px)'
footerClassName='grid grid-cols-2 gap-2 sm:flex' footerClassName='grid grid-cols-2 gap-2 sm:flex'
......
...@@ -77,6 +77,12 @@ const LazyUserCharts = lazy(() => ...@@ -77,6 +77,12 @@ const LazyUserCharts = lazy(() =>
})) }))
) )
const LazyFlowCharts = lazy(() =>
import('./components/flow/flow-charts').then((m) => ({
default: m.FlowCharts,
}))
)
function LogStatCardsFallback() { function LogStatCardsFallback() {
return ( return (
<div className='overflow-hidden rounded-lg border'> <div className='overflow-hidden rounded-lg border'>
...@@ -137,6 +143,9 @@ const SECTION_META: Record<DashboardSectionId, { titleKey: string }> = { ...@@ -137,6 +143,9 @@ const SECTION_META: Record<DashboardSectionId, { titleKey: string }> = {
models: { models: {
titleKey: 'Model Call Analytics', titleKey: 'Model Call Analytics',
}, },
flow: {
titleKey: 'Flow',
},
users: { users: {
titleKey: 'User Analytics', titleKey: 'User Analytics',
}, },
...@@ -217,6 +226,17 @@ export function Dashboard() { ...@@ -217,6 +226,17 @@ export function Dashboard() {
/> />
</> </>
) : null ) : null
const flowActions =
activeSection === 'flow' ? (
<ModelsFilter
preferences={chartPreferences}
onFilterChange={handleFilterChange}
onReset={handleResetFilters}
titleKey='Flow Filters'
descriptionKey='Filter the traffic flow view by time range and user.'
/>
) : null
const sectionActions = modelActions ?? flowActions
return ( return (
<SectionPageLayout> <SectionPageLayout>
...@@ -238,9 +258,9 @@ export function Dashboard() { ...@@ -238,9 +258,9 @@ export function Dashboard() {
) : ( ) : (
<div /> <div />
)} )}
{modelActions != null && ( {sectionActions != null && (
<div className='flex shrink-0 flex-wrap items-center gap-1.5 sm:gap-2'> <div className='flex shrink-0 flex-wrap items-center gap-1.5 sm:gap-2'>
{modelActions} {sectionActions}
</div> </div>
)} )}
</div> </div>
...@@ -298,6 +318,13 @@ export function Dashboard() { ...@@ -298,6 +318,13 @@ export function Dashboard() {
</Suspense> </Suspense>
</FadeIn> </FadeIn>
)} )}
{activeSection === 'flow' && (
<FadeIn>
<Suspense fallback={<ModelChartsFallback />}>
<LazyFlowCharts filters={modelFilters} />
</Suspense>
</FadeIn>
)}
</div> </div>
</SectionPageLayout.Content> </SectionPageLayout.Content>
</SectionPageLayout> </SectionPageLayout>
......
...@@ -38,13 +38,15 @@ type TooltipLineItem = { ...@@ -38,13 +38,15 @@ type TooltipLineItem = {
shapeSize?: number shapeSize?: number
} }
function getVChartDefaultColors(domainLength: number) { export function getDashboardChartColors(domainLength: number): string[] {
const scheme = const scheme =
vchartDefaultDataScheme.find( vchartDefaultDataScheme.find(
(item) => !item.maxDomainLength || domainLength <= item.maxDomainLength (item) => !item.maxDomainLength || domainLength <= item.maxDomainLength
) ?? vchartDefaultDataScheme[vchartDefaultDataScheme.length - 1] ) ?? vchartDefaultDataScheme[vchartDefaultDataScheme.length - 1]
return scheme.scheme return scheme.scheme.filter(
(color): color is string => typeof color === 'string'
)
} }
function renderQuotaCompat(rawQuota: number, digits = 4): string { function renderQuotaCompat(rawQuota: number, digits = 4): string {
...@@ -259,7 +261,7 @@ export function processChartData( ...@@ -259,7 +261,7 @@ export function processChartData(
const sortedTimes = Array.from(timeModelMap.keys()).sort() const sortedTimes = Array.from(timeModelMap.keys()).sort()
const sortedModels = [...allModels].sort() const sortedModels = [...allModels].sort()
const modelColorDomain = Array.from(new Set([...sortedModels, otherLabel])) const modelColorDomain = Array.from(new Set([...sortedModels, otherLabel]))
const modelColorRange = getVChartDefaultColors(modelColorDomain.length) const modelColorRange = getDashboardChartColors(modelColorDomain.length)
const otherColor = modelColorRange[modelColorDomain.indexOf(otherLabel)] const otherColor = modelColorRange[modelColorDomain.indexOf(otherLabel)]
const otherTooltipColor = const otherTooltipColor =
typeof otherColor === 'string' ? otherColor : '#FF8A00' typeof otherColor === 'string' ? otherColor : '#FF8A00'
......
import assert from 'node:assert/strict'
import { describe, test } from 'node:test'
import type { FlowUserFilterOption } from '../types'
import {
compactFlowSelectionLabel,
flowDisplayState,
requireSuccessfulFlowRows,
visibleFlowUsers,
} from './flow-selection'
const users: FlowUserFilterOption[] = [
{
value: 'user:1',
label: 'dry',
valueLabel: '100',
valueRaw: 100,
color: '#1664ff',
},
{
value: 'user:2',
label: 'jrc',
valueLabel: '70',
valueRaw: 70,
color: '#1ac6ff',
},
]
describe('dashboard flow selection helpers', () => {
test('limits user chips to currently visible users', () => {
assert.deepEqual(
visibleFlowUsers(users, []).map((user) => user.value),
['user:1', 'user:2']
)
assert.deepEqual(
visibleFlowUsers(users, ['user:2']).map((user) => user.value),
['user:2']
)
})
test('filters visible users without mutating the source options', () => {
const visible = visibleFlowUsers(users, ['user:1'])
assert.deepEqual(
visible.map((user) => user.value),
['user:1']
)
assert.deepEqual(
users.map((user) => user.value),
['user:1', 'user:2']
)
})
test('formats compact selected counts for flow multiselect summaries', () => {
assert.equal(compactFlowSelectionLabel(0), '*')
assert.equal(compactFlowSelectionLabel(1), '1')
assert.equal(compactFlowSelectionLabel(23), '23')
})
test('prioritizes loading and error states before empty flow data', () => {
assert.equal(
flowDisplayState({
isLoading: true,
isError: true,
linkCount: 0,
themeReady: true,
}),
'loading'
)
assert.equal(
flowDisplayState({
isLoading: false,
isError: true,
linkCount: 0,
themeReady: true,
}),
'error'
)
assert.equal(
flowDisplayState({
isLoading: false,
isError: false,
linkCount: 0,
themeReady: true,
}),
'empty'
)
assert.equal(
flowDisplayState({
isLoading: false,
isError: false,
linkCount: 1,
themeReady: false,
}),
'loading'
)
})
test('throws unsuccessful flow responses instead of treating them as empty data', () => {
assert.throws(
() =>
requireSuccessfulFlowRows(
{ success: false, data: [], message: 'database unavailable' },
'Failed to load'
),
/database unavailable/
)
assert.deepEqual(
requireSuccessfulFlowRows(
{ success: true, data: [{ user_id: 1, quota: 10 }] },
'Failed to load'
),
[{ user_id: 1, quota: 10 }]
)
})
})
/*
Copyright (C) 2023-2026 QuantumNous
This program is free software: you can redistribute it and/or modify
it under the terms of the GNU Affero General Public License as
published by the Free Software Foundation, either version 3 of the
License, or (at your option) any later version.
This program is distributed in the hope that it will be useful,
but WITHOUT ANY WARRANTY; without even the implied warranty of
MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
GNU Affero General Public License for more details.
You should have received a copy of the GNU Affero General Public License
along with this program. If not, see <https://www.gnu.org/licenses/>.
For commercial licensing, please contact support@quantumnous.com
*/
import type {
FlowQuotaDataItem,
FlowUserFilterOption,
} from '@/features/dashboard/types'
export type FlowDisplayState = 'loading' | 'error' | 'empty' | 'chart'
export interface FlowResponse {
success: boolean
data?: FlowQuotaDataItem[]
message?: string
}
export function requireSuccessfulFlowRows(
response: FlowResponse,
fallbackMessage: string
): FlowQuotaDataItem[] {
if (!response.success) {
throw new Error(response.message || fallbackMessage)
}
return response.data ?? []
}
export function flowDisplayState(options: {
isLoading: boolean
isError: boolean
linkCount: number
themeReady: boolean
}): FlowDisplayState {
if (options.isLoading) return 'loading'
if (options.isError) return 'error'
if (options.linkCount === 0) return 'empty'
if (!options.themeReady) return 'loading'
return 'chart'
}
export function compactFlowSelectionLabel(count: number): string {
return count > 0 ? String(count) : '*'
}
export function visibleFlowUsers(
users: FlowUserFilterOption[],
selectedUsers: string[]
): FlowUserFilterOption[] {
if (selectedUsers.length === 0) return users
const selected = new Set(selectedUsers)
return users.filter((user) => selected.has(user.value))
}
import assert from 'node:assert/strict'
import { describe, test } from 'node:test'
import type { FlowQuotaDataItem } from '../types'
import {
buildDashboardFlowData,
buildFlowFilterOptions,
buildFlowSankeySpec,
} from './flow'
const rows: FlowQuotaDataItem[] = [
{
user_id: 1,
username: 'alice',
node_name: 'node-a',
token_id: 11,
token_name: 'primary',
use_group: 'vip',
channel_id: 101,
channel_name: 'east',
model_name: 'gpt-4.1',
quota: 100,
token_used: 40,
count: 2,
},
{
user_id: 1,
username: 'alice',
node_name: 'node-a',
token_id: 11,
token_name: 'primary',
use_group: 'vip',
channel_id: 102,
channel_name: 'west',
model_name: 'gpt-4.1',
quota: 50,
token_used: 20,
count: 1,
},
{
user_id: 2,
username: 'bob',
node_name: 'node-b',
token_id: 22,
token_name: 'backup',
use_group: 'default',
channel_id: 101,
channel_name: 'east',
model_name: 'claude-4-sonnet',
quota: 70,
token_used: 30,
count: 3,
},
]
describe('dashboard flow data', () => {
test('builds normal user token-group-model flow', () => {
const result = buildDashboardFlowData(rows.slice(0, 2), 'quota', {
role: 'user',
})
assert.equal(result.summary.quota, 150)
assert.equal(result.summary.tokens, 60)
assert.equal(result.summary.requests, 3)
assert.deepEqual(
result.flow.links.map((link) => [link.source, link.target, link.value]),
[
['group:vip', 'model:gpt-4.1', 150],
['token:11', 'group:vip', 150],
]
)
assert.equal(
result.flow.nodes.some((node) => node.kind === 'channel'),
false
)
})
test('builds admin user-group-model-channel flow', () => {
const result = buildDashboardFlowData(rows, 'quota', {
role: 'admin',
})
assert.deepEqual(
result.flow.links.map((link) => [link.source, link.target, link.value]),
[
['group:default', 'model:claude-4-sonnet', 70],
['group:vip', 'model:gpt-4.1', 150],
['model:claude-4-sonnet', 'channel:101', 70],
['model:gpt-4.1', 'channel:101', 100],
['model:gpt-4.1', 'channel:102', 50],
['user:1', 'group:vip', 150],
['user:2', 'group:default', 70],
]
)
})
test('builds root user-node-token-group-model-channel flow', () => {
const result = buildDashboardFlowData(rows, 'requests', {
role: 'root',
})
assert.deepEqual(
result.flow.links.map((link) => [link.source, link.target, link.value]),
[
['group:default', 'model:claude-4-sonnet', 3],
['group:vip', 'model:gpt-4.1', 3],
['model:claude-4-sonnet', 'channel:101', 3],
['model:gpt-4.1', 'channel:101', 2],
['model:gpt-4.1', 'channel:102', 1],
['node:node-a', 'token:11', 3],
['node:node-b', 'token:22', 3],
['token:11', 'group:vip', 3],
['token:22', 'group:default', 3],
['user:1', 'node:node-a', 3],
['user:2', 'node:node-b', 3],
]
)
})
test('filters by selected users', () => {
const result = buildDashboardFlowData(rows, 'quota', {
role: 'admin',
selectedUsers: ['user:2'],
})
assert.equal(result.summary.quota, 70)
assert.deepEqual(
result.flow.links.map((link) => [link.source, link.target, link.value]),
[
['group:default', 'model:claude-4-sonnet', 70],
['model:claude-4-sonnet', 'channel:101', 70],
['user:2', 'group:default', 70],
]
)
})
test('reconnects links when a middle stage is hidden', () => {
const result = buildDashboardFlowData(rows, 'quota', {
role: 'admin',
visibleStages: ['user', 'model', 'channel'],
})
assert.deepEqual(
result.flow.links.map((link) => [link.source, link.target, link.value]),
[
['model:claude-4-sonnet', 'channel:101', 70],
['model:gpt-4.1', 'channel:101', 100],
['model:gpt-4.1', 'channel:102', 50],
['user:1', 'model:gpt-4.1', 150],
['user:2', 'model:claude-4-sonnet', 70],
]
)
assert.equal(
result.flow.nodes.some((node) => node.kind === 'group'),
false
)
})
test('ignores stage filters that would leave fewer than two columns', () => {
const result = buildDashboardFlowData(rows.slice(0, 2), 'quota', {
role: 'user',
visibleStages: ['model'],
})
assert.deepEqual(
result.flow.links.map((link) => [link.source, link.target, link.value]),
[
['group:vip', 'model:gpt-4.1', 150],
['token:11', 'group:vip', 150],
]
)
})
test('builds user filter options with stable values', () => {
const options = buildFlowFilterOptions(rows, 'quota')
assert.deepEqual(
options.users.map((user) => [user.value, user.label, user.valueLabel]),
[
['user:1', 'alice', '150'],
['user:2', 'bob', '70'],
]
)
assert.notEqual(options.users[0].color, options.users[1].color)
})
test('builds Sankey spec with quota token request tooltips', () => {
const result = buildDashboardFlowData(rows.slice(0, 1), 'quota', {
role: 'root',
})
const flowSpec = buildFlowSankeySpec(result.flow, 'Flow')
const values = flowSpec.data[0].values[0]
const aliceNode = values.nodes.find(
(node: Record<string, unknown>) => node.key === 'user:1'
)
const userNodeLink = values.links.find(
(link: Record<string, unknown>) =>
link.source === 'user:1' && link.target === 'node:node-a'
)
assert.equal(flowSpec.type, 'sankey')
assert.equal(flowSpec.title.text, 'Flow')
assert.equal(flowSpec.tooltip.mark.visible({ datum: aliceNode }), true)
assert.equal(flowSpec.tooltip.mark.visible({ datum: userNodeLink }), true)
assert.equal(flowSpec.animation, false)
assert.equal(values.nodes.length, 6)
assert.equal(values.links.length, 5)
assert.equal(aliceNode.name, 'alice')
assert.match(userNodeLink.linkColor, /^rgba\(/)
const tooltipRows = flowSpec.tooltip.mark.content
assert.deepEqual(
tooltipRows
.filter((row: Record<string, unknown>) =>
typeof row.visible === 'function'
? row.visible({ datum: userNodeLink })
: true
)
.map((row: Record<string, unknown>) => [
row.key,
typeof row.value === 'function'
? row.value({ datum: userNodeLink })
: row.value,
]),
[
['Quota', '100'],
['Tokens', '40'],
['Requests', '2'],
['Share', '100.0%'],
]
)
})
})
...@@ -33,5 +33,10 @@ export { ...@@ -33,5 +33,10 @@ export {
getDefaultPingStatus, getDefaultPingStatus,
} from './api-info' } from './api-info'
export { processChartData, processUserChartData } from './charts' export { processChartData, processUserChartData } from './charts'
export {
buildDashboardFlowData,
buildFlowSankeySpec,
getFlowStages,
} from './flow'
export { safeDivide, calculateDashboardStats } from './stats' export { safeDivide, calculateDashboardStats } from './stats'
export { getPreviewText } from './text' export { getPreviewText } from './text'
...@@ -34,6 +34,11 @@ const DASHBOARD_SECTIONS = [ ...@@ -34,6 +34,11 @@ const DASHBOARD_SECTIONS = [
build: () => null, build: () => null,
}, },
{ {
id: 'flow',
titleKey: 'Flow',
build: () => null,
},
{
id: 'users', id: 'users',
titleKey: 'User Analytics', titleKey: 'User Analytics',
adminOnly: true, adminOnly: true,
......
...@@ -33,6 +33,101 @@ export interface QuotaDataItem { ...@@ -33,6 +33,101 @@ export interface QuotaDataItem {
quota?: number quota?: number
} }
export interface FlowQuotaDataItem {
user_id?: number
username?: string
node_name?: string
use_group?: string
token_id?: number
token_name?: string
channel_id?: number
channel_name?: string
model_name?: string
token_used?: number
count?: number
quota?: number
}
export type FlowMetric = 'quota' | 'tokens' | 'requests'
export type FlowRole = 'user' | 'admin' | 'root'
export type FlowNodeKind =
| 'user'
| 'node'
| 'token'
| 'group'
| 'model'
| 'channel'
export interface FlowBuildOptions {
role?: FlowRole
selectedUsers?: string[]
colorPalette?: readonly string[]
visibleStages?: FlowNodeKind[]
// Resolves the label for a token whose record no longer exists (deleted).
// Lets the caller inject a localized string such as "Deleted (123)".
deletedTokenLabel?: (tokenId: number) => string
}
export interface DashboardFlowNode {
id: string
label: string
kind: FlowNodeKind
value: number
requests: number
quota: number
tokens: number
color: string
colorKey: string
}
export interface DashboardFlowLink {
source: string
target: string
value: number
requests: number
quota: number
tokens: number
sourceLabel: string
targetLabel: string
color: string
linkColor: string
linkAlpha: number
hoverColor: string
colorKey: string
share: number
}
export interface DashboardFlowGraph {
nodes: DashboardFlowNode[]
links: DashboardFlowLink[]
}
export interface FlowUserFilterOption {
value: string
label: string
valueLabel: string
valueRaw: number
color: string
}
export interface FlowFilterOptions {
users: FlowUserFilterOption[]
}
export interface FlowSummary {
quota: number
tokens: number
requests: number
}
export interface ProcessedFlowData {
summary: FlowSummary
flow: DashboardFlowGraph
filterOptions: FlowFilterOptions
}
// ============================================================================ // ============================================================================
// Uptime Monitoring Types // Uptime Monitoring Types
// ============================================================================ // ============================================================================
......
...@@ -17,13 +17,13 @@ ...@@ -17,13 +17,13 @@
"file": "ja.json", "file": "ja.json",
"missingCount": 0, "missingCount": 0,
"extrasCount": 0, "extrasCount": 0,
"untranslatedCount": 1 "untranslatedCount": 0
}, },
"ru": { "ru": {
"file": "ru.json", "file": "ru.json",
"missingCount": 0, "missingCount": 0,
"extrasCount": 0, "extrasCount": 0,
"untranslatedCount": 1 "untranslatedCount": 0
}, },
"vi": { "vi": {
"file": "vi.json", "file": "vi.json",
......
...@@ -435,6 +435,10 @@ export const STATIC_I18N_KEYS = [ ...@@ -435,6 +435,10 @@ export const STATIC_I18N_KEYS = [
'Data management and log viewing', 'Data management and log viewing',
'Dashboard', 'Dashboard',
'System data statistics', 'System data statistics',
'Flow',
'Flow Filters',
'Filter the traffic flow view by time range and user.',
'Requests',
'Token Management', 'Token Management',
'API token management', 'API token management',
'Usage Logs', 'Usage Logs',
...@@ -490,6 +494,20 @@ export const STATIC_I18N_KEYS = [ ...@@ -490,6 +494,20 @@ export const STATIC_I18N_KEYS = [
'Batch detection failed', 'Batch detection failed',
'Batch detection complete: {{channels}} channels, {{add}} to add, {{remove}} to remove, {{fails}} failed', 'Batch detection complete: {{channels}} channels, {{add}} to add, {{remove}} to remove, {{fails}} failed',
// Dashboard flow stages (labels/descriptions passed to t at runtime)
'User',
'Node',
'Token',
'Group',
'Model',
'Channel',
'The user who made the requests',
'The deployment node that handled the requests',
'The API key used for the requests',
'The user group applied to the requests',
'The model that was requested',
'The upstream channel that served the requests',
// Misc // Misc
'Cancel', 'Cancel',
'Status', 'Status',
......
Markdown is supported
0% or
You are about to add 0 people to the discussion. Proceed with caution.
Finish editing this message first!
Please register or sign in to comment