cluster

package
v1.0.44 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Jun 26, 2026 License: MIT Imports: 46 Imported by: 0

Documentation

Index

Constants

View Source
const (
	BillingBalanceTypeRecharge = "recharge"
	BillingBalanceTypeDeduct   = "deduct"

	BillingPriceSourceManual  = "manual"
	BillingPriceSourceDefault = "default"
	BillingPriceSourceSync    = "sync"
)
View Source
const (
	ProxyPoolScopeGlobal = "global"

	ProxyPoolTestResultPassed   = "passed"
	ProxyPoolTestResultFailed   = "failed"
	ProxyPoolTestResultUntested = "untested"
)
View Source
const DefaultConfigPath = "cluster.yaml"
View Source
const OAuthSessionTTL = 20 * time.Minute

Variables

View Source
var ErrAPIKeyExists = errors.New("api key already exists")

ErrAPIKeyExists indicates that an API key value is already owned by another record.

View Source
var ErrBillingDuplicateModelPrice = errors.New("billing model price already exists")
View Source
var ErrUserNotFound = errors.New("user not found")

ErrUserNotFound indicates that the referenced user record does not exist.

Functions

func APIKeyHash

func APIKeyHash(apiKey string) string

APIKeyHash handles an api key hash.

func ApplyOriginalAuthFileName

func ApplyOriginalAuthFileName(auths []*coreauth.Auth, filename string)

ApplyOriginalAuthFileName stores the source OAuth file name for management display.

func AutoMigrate

func AutoMigrate(db *gorm.DB) error

AutoMigrate handles an auto migrate.

func ClientAddr

func ClientAddr(ctx context.Context, db *gorm.DB) (string, error)

ClientAddr handles a client addr.

func ConfigRootFromSnapshot

func ConfigRootFromSnapshot(snapshot map[string]json.RawMessage) (map[string]any, error)

ConfigRootFromSnapshot derives config root from snapshot.

func DeterministicAPIKeyUUID

func DeterministicAPIKeyUUID(provider, baseURL, apiKeyHash, compatName, providerKey string) string

DeterministicAPIKeyUUID handles a deterministic api key uuid.

func EnsureOAuthPayloadUUID

func EnsureOAuthPayloadUUID(raw []byte) ([]byte, string, bool, error)

EnsureOAuthPayloadUUID validates ensure o auth payload uuid.

func ForwardRefreshToMaster

func ForwardRefreshToMaster(ctx context.Context, master *ClusterNodeRecord, authUUID string, secret string, tlsConfig *tls.Config) ([]byte, error)

ForwardRefreshToMaster converts forward refresh to master.

func HTTPClientTLSConfig added in v1.0.44

func HTTPClientTLSConfig(tlsConfig *tls.Config, serverName string) (*tls.Config, error)

HTTPClientTLSConfig derives an mTLS client config for cluster HTTP calls.

func IsUserConflictError added in v1.0.44

func IsUserConflictError(err error) bool

IsUserConflictError reports whether an error is a username uniqueness conflict.

func OAuthSessionData

func OAuthSessionData(record *OAuthSessionRecord) (map[string]any, error)

OAuthSessionData handles an o auth session data.

func Open

func Open(ctx context.Context, cfg PGSQLConfig) (*gorm.DB, error)

Open opens the resource.

func OpenSQLite added in v1.0.44

func OpenSQLite(ctx context.Context, path string) (*gorm.DB, error)

OpenSQLite opens a SQLite database.

func ParseChannelRecordID added in v1.0.44

func ParseChannelRecordID(value string) (uint, error)

ParseChannelRecordID parses a positive unsigned record ID.

func ParseModelRecordID added in v1.0.44

func ParseModelRecordID(value string) (uint, error)

ParseModelRecordID parses a positive unsigned record ID.

func ParseUserRecordID added in v1.0.44

func ParseUserRecordID(value string) (uint, error)

ParseUserRecordID parses a positive unsigned user record ID.

func RecordToAuth

func RecordToAuth(record *AuthRecord) (*coreauth.Auth, error)

RecordToAuth records a to auth.

func RuntimeConfigFromRoot

func RuntimeConfigFromRoot(root map[string]any) (*appconfig.Config, []byte, error)

RuntimeConfigFromRoot derives runtime config from root.

Types

type APIKeyEntry added in v1.0.44

type APIKeyEntry struct {
	APIKey      string
	UserID      *uint
	Channels    []uint
	ModelGroups []uint
}

func APIKeyEntryFromRecord added in v1.0.44

func APIKeyEntryFromRecord(record *APIKeyRecord) (APIKeyEntry, error)

APIKeyEntryFromRecord converts an API key record to a response entry.

type APIKeyEntryUpdate added in v1.0.44

type APIKeyEntryUpdate struct {
	APIKey      string
	UserID      *uint
	Channels    *[]uint
	ModelGroups *[]uint
}

type APIKeyRecord

type APIKeyRecord struct {
	ID          uint           `gorm:"column:id;primaryKey;autoIncrement;index:idx_api_key_active_order,priority:2"`
	APIKey      string         `gorm:"column:api_key;not null;uniqueIndex"`
	UserID      *uint          `gorm:"column:user_id;index;index:idx_api_key_user_active,priority:1"`
	Channels    JSONB          `gorm:"column:channels"`
	ModelGroups JSONB          `gorm:"column:model_groups"`
	CreatedAt   time.Time      `gorm:"column:created_at"`
	UpdatedAt   time.Time      `gorm:"column:updated_at"`
	DeletedAt   gorm.DeletedAt `gorm:"column:deleted_at;index;index:idx_api_key_active_order,priority:1;index:idx_api_key_user_active,priority:2"`
}

func (APIKeyRecord) TableName

func (APIKeyRecord) TableName() string

TableName returns the database table name.

type APIKeyUpsertStats added in v1.0.44

type APIKeyUpsertStats struct {
	Created   int
	Updated   int
	Unchanged int
	Restored  int
	Removed   int
}

func (APIKeyUpsertStats) Changed added in v1.0.44

func (s APIKeyUpsertStats) Changed() bool

Changed reports whether api key rows were mutated.

type APIKeyUserUpdate added in v1.0.44

type APIKeyUserUpdate struct {
	APIKey      *string
	Channels    *[]uint
	ModelGroups *[]uint
}

type AppLogQuery added in v1.0.44

type AppLogQuery struct {
	HomeIP    string
	ClientIP  string
	RequestID string
	Level     string
	After     *time.Time
	Before    *time.Time
	Limit     int
	Offset    int
}

type AppLogQueryResult added in v1.0.44

type AppLogQueryResult struct {
	Records []AppLogRecord
	Total   int64
}

type AppLogRecord added in v1.0.44

type AppLogRecord struct {
	ID        uint      `gorm:"column:id;primaryKey;autoIncrement;index:idx_log_time_order,priority:2"`
	Timestamp time.Time `` /* 195-byte string literal not displayed */
	ClientIP  string    `gorm:"column:client_ip;index:idx_log_client_ip;index:idx_log_client_time,priority:1"`
	RequestID string    `gorm:"column:request_id;index:idx_log_request_id;index:idx_log_home_request,priority:2"`
	HomeIP    string    `gorm:"column:home_ip;index:idx_log_home_ip;index:idx_log_home_request,priority:1"`
	Level     string    `gorm:"column:level;index:idx_log_level;index:idx_log_level_time,priority:1"`
	Line      string    `gorm:"column:line;type:text;not null"`
	CreatedAt time.Time `gorm:"column:created_at;not null;index:idx_log_created_at"`
}

func AppLogRecordFromPayload added in v1.0.44

func AppLogRecordFromPayload(clientIP string, homeIP string, payload string) (*AppLogRecord, error)

AppLogRecordFromPayload creates an app log record from a CPA payload.

func (AppLogRecord) TableName added in v1.0.44

func (AppLogRecord) TableName() string

TableName returns the database table name.

type AuthIndex

type AuthIndex struct {
	UUID          string
	ID            string
	Index         string
	Provider      string
	Label         string
	Prefix        string
	Status        coreauth.Status
	Disabled      bool
	Unavailable   bool
	BaseURL       string
	ModelsHash    string
	Attributes    map[string]string
	ModelMetadata map[string]any
}

type AuthRecord

type AuthRecord struct {
	UUID             string          `gorm:"column:uuid;primaryKey"`
	AuthJSON         JSONB           `gorm:"column:auth_json;not null"`
	Version          int64           `gorm:"column:version;not null;default:1"`
	ID               string          `gorm:"column:id;index:idx_auth_active_order,priority:2"`
	Index            string          `gorm:"column:index"`
	Provider         string          `gorm:"column:provider"`
	Label            string          `gorm:"column:label"`
	Prefix           string          `gorm:"column:prefix"`
	Status           coreauth.Status `gorm:"column:status"`
	Disabled         bool            `gorm:"column:disabled"`
	Unavailable      bool            `gorm:"column:unavailable"`
	BaseURL          string          `gorm:"column:base_url"`
	APIKeyHash       string          `gorm:"column:api_key_hash"`
	CompatName       string          `gorm:"column:compat_name"`
	ProviderKey      string          `gorm:"column:provider_key"`
	ModelsHash       string          `gorm:"column:models_hash"`
	CreatedAt        time.Time       `gorm:"column:created_at"`
	UpdatedAt        time.Time       `gorm:"column:updated_at"`
	LastRefreshedAt  *time.Time      `gorm:"column:last_refreshed_at"`
	NextRefreshAfter *time.Time      `gorm:"column:next_refresh_after"`
	DeletedAt        gorm.DeletedAt  `gorm:"column:deleted_at;index;index:idx_auth_active_order,priority:1"`
}

func AuthToRecord

func AuthToRecord(auth *coreauth.Auth) (*AuthRecord, error)

AuthToRecord converts auth to record.

func (AuthRecord) TableName

func (AuthRecord) TableName() string

TableName returns the database table name.

type BillingBalanceQuery added in v1.0.44

type BillingBalanceQuery struct {
	From     *time.Time
	To       *time.Time
	UserText string
	UserID   *uint
	Limit    int
	Offset   int
}

type BillingBalanceRecord added in v1.0.44

type BillingBalanceRecord struct {
	ID            string    `gorm:"column:id;primaryKey"`
	UserID        uint      `gorm:"column:user_id;not null;index:idx_billing_balance_user_time,priority:1"`
	Type          string    `gorm:"column:type;not null;index:idx_billing_balance_type_time,priority:1"`
	Amount        float64   `gorm:"column:amount;not null"`
	BalanceBefore float64   `gorm:"column:balance_before;not null"`
	BalanceAfter  float64   `gorm:"column:balance_after;not null"`
	Operator      string    `gorm:"column:operator;not null"`
	Note          string    `gorm:"column:note;type:text"`
	CreatedAt     time.Time `` /* 147-byte string literal not displayed */
}

func (BillingBalanceRecord) TableName added in v1.0.44

func (BillingBalanceRecord) TableName() string

type BillingBalanceResult added in v1.0.44

type BillingBalanceResult struct {
	Records []BillingBalanceRecord
	Total   int64
}

type BillingBalanceUpdate added in v1.0.44

type BillingBalanceUpdate struct {
	UserID   uint
	Type     string
	Amount   float64
	Operator string
	Note     string
}

type BillingChargeQuery added in v1.0.44

type BillingChargeQuery struct {
	From     *time.Time
	To       *time.Time
	UserText string
	UserID   *uint
	Provider string
	Model    string
	Limit    int
	Offset   int
}

type BillingChargeRecord added in v1.0.44

type BillingChargeRecord struct {
	ID               string    `gorm:"column:id;primaryKey"`
	UsageID          uint      `gorm:"column:usage_id;not null;uniqueIndex"`
	PayloadHash      string    `gorm:"column:payload_hash;not null;uniqueIndex"`
	UserID           *uint     `gorm:"column:user_id;index:idx_billing_charge_user_time,priority:1"`
	APIKeyID         *uint     `gorm:"column:api_key_id;index"`
	APIKeyLabel      string    `gorm:"column:api_key_label"`
	APIKeyMasked     string    `gorm:"column:api_key_masked"`
	Provider         string    `gorm:"column:provider;index:idx_billing_charge_provider_model_time,priority:1"`
	Model            string    `gorm:"column:model;index:idx_billing_charge_provider_model_time,priority:2"`
	OriginalModel    string    `gorm:"column:original_model"`
	ActualModel      string    `gorm:"column:actual_model"`
	RequestID        string    `gorm:"column:request_id;index"`
	Endpoint         string    `gorm:"column:endpoint;index"`
	InputTokens      int64     `gorm:"column:input_tokens;not null;default:0"`
	OutputTokens     int64     `gorm:"column:output_tokens;not null;default:0"`
	CacheTokens      int64     `gorm:"column:cache_tokens;not null;default:0"`
	Amount           float64   `gorm:"column:amount;not null;default:0"`
	BalanceBefore    float64   `gorm:"column:balance_before;not null;default:0"`
	BalanceAfter     float64   `gorm:"column:balance_after;not null;default:0"`
	MatchedPriceRule string    `gorm:"column:matched_price_rule"`
	PriceSnapshot    JSONB     `gorm:"column:price_snapshot;not null"`
	RequestSummary   string    `gorm:"column:request_summary;type:text"`
	CreatedAt        time.Time `` /* 155-byte string literal not displayed */
}

func (BillingChargeRecord) TableName added in v1.0.44

func (BillingChargeRecord) TableName() string

type BillingChargeResult added in v1.0.44

type BillingChargeResult struct {
	Records []BillingChargeRecord
	Total   int64
}

type BillingModelPricePatch added in v1.0.44

type BillingModelPricePatch struct {
	Provider                  *string
	Model                     *string
	InputPricePerMillion      *float64
	OutputPricePerMillion     *float64
	CacheReadPricePerMillion  *float64
	CacheWritePricePerMillion *float64
	RequestPrice              *float64
	Source                    *string
	Enabled                   *bool
	Note                      *string
}

type BillingModelPriceQuery added in v1.0.44

type BillingModelPriceQuery struct {
	Provider string
	Model    string
	Enabled  *bool
}

type BillingModelPriceRecord added in v1.0.44

type BillingModelPriceRecord struct {
	ID                        string         `gorm:"column:id;primaryKey"`
	Provider                  string         `gorm:"column:provider;not null;index:idx_billing_model_price_lookup,priority:1"`
	Model                     string         `gorm:"column:model;not null;index:idx_billing_model_price_lookup,priority:2"`
	InputPricePerMillion      float64        `gorm:"column:input_price_per_million;not null;default:0"`
	OutputPricePerMillion     float64        `gorm:"column:output_price_per_million;not null;default:0"`
	CacheReadPricePerMillion  float64        `gorm:"column:cache_read_price_per_million;not null;default:0"`
	CacheWritePricePerMillion float64        `gorm:"column:cache_write_price_per_million;not null;default:0"`
	RequestPrice              float64        `gorm:"column:request_price;not null;default:0"`
	Source                    string         `gorm:"column:source;not null;default:manual"`
	Enabled                   bool           `gorm:"column:enabled;not null;index:idx_billing_model_price_enabled"`
	Note                      string         `gorm:"column:note;type:text"`
	CreatedAt                 time.Time      `gorm:"column:created_at"`
	UpdatedAt                 time.Time      `gorm:"column:updated_at"`
	DeletedAt                 gorm.DeletedAt `gorm:"column:deleted_at;index"`
}

func (BillingModelPriceRecord) TableName added in v1.0.44

func (BillingModelPriceRecord) TableName() string

type BillingModelPriceUpdate added in v1.0.44

type BillingModelPriceUpdate struct {
	Provider                  string
	Model                     string
	InputPricePerMillion      float64
	OutputPricePerMillion     float64
	CacheReadPricePerMillion  float64
	CacheWritePricePerMillion float64
	RequestPrice              float64
	Source                    string
	Enabled                   bool
	Note                      string
}

type BillingOverview added in v1.0.44

type BillingOverview struct {
	Range               BillingRange
	TotalChargeAmount   float64
	TotalRechargeAmount float64
	TotalDeductAmount   float64
	TotalBalance        float64
	RequestCount        int64
	InputTokens         int64
	OutputTokens        int64
	CacheTokens         int64
	ActiveUserCount     int64
	DailyTrend          []BillingTrendPoint
	TopUsers            []BillingTopItem
	TopModels           []BillingTopItem
	TopProviders        []BillingTopItem
}

type BillingOverviewQuery added in v1.0.44

type BillingOverviewQuery struct {
	From     *time.Time
	To       *time.Time
	UserText string
	UserID   *uint
	Provider string
	Model    string
}

type BillingPriceSnapshot added in v1.0.44

type BillingPriceSnapshot struct {
	Provider                  string  `json:"provider"`
	Model                     string  `json:"model"`
	InputPricePerMillion      float64 `json:"input_price_per_million"`
	OutputPricePerMillion     float64 `json:"output_price_per_million"`
	CacheReadPricePerMillion  float64 `json:"cache_read_price_per_million"`
	CacheWritePricePerMillion float64 `json:"cache_write_price_per_million"`
	RequestPrice              float64 `json:"request_price"`
}

type BillingRange added in v1.0.44

type BillingRange struct {
	From string
	To   string
}

type BillingTopItem added in v1.0.44

type BillingTopItem struct {
	ID           string
	Label        string
	Amount       float64
	RequestCount int64
}

type BillingTrendPoint added in v1.0.44

type BillingTrendPoint struct {
	Date         string
	ChargeAmount float64
	RequestCount int64
}

type CertificateRecord added in v1.0.44

type CertificateRecord struct {
	ID                     string    `gorm:"column:id;primaryKey;index:idx_certificate_ca_order,priority:3"`
	ClusterID              string    `gorm:"column:cluster_id;index"`
	CertificatePEM         string    `gorm:"column:certificate_pem;type:text"`
	CertificateFingerprint string    `gorm:"column:certificate_fingerprint;index;index:idx_certificate_peer_lookup,priority:1"`
	PrivateKeyPEM          string    `gorm:"column:private_key_pem;type:text"`
	CSRPEM                 string    `gorm:"column:csr_pem;type:text"`
	IP                     string    `gorm:"column:ip;index:idx_certificate_server_ip,priority:2"`
	CAFingerprint          string    `gorm:"column:ca_fingerprint"`
	EnrollmentSecretHash   string    `gorm:"column:enrollment_secret_hash"`
	IsCA                   bool      `gorm:"column:is_ca;index;index:idx_certificate_ca_order,priority:1"`
	IsServer               bool      `gorm:"column:is_server;index:idx_certificate_server_ip,priority:1;index:idx_certificate_peer_lookup,priority:3"`
	IsClient               bool      `gorm:"column:is_client;index;index:idx_certificate_peer_lookup,priority:2"`
	SerialNumber           string    `gorm:"column:serial_number"`
	NotBefore              time.Time `gorm:"column:not_before"`
	NotAfter               time.Time `gorm:"column:not_after"`
	CreatedAt              time.Time `gorm:"column:created_at;index:idx_certificate_ca_order,priority:2"`
	UpdatedAt              time.Time `gorm:"column:updated_at"`
}

func (CertificateRecord) TableName added in v1.0.44

func (CertificateRecord) TableName() string

TableName returns the database table name.

type ChannelGroupDetailFilter added in v1.0.44

type ChannelGroupDetailFilter struct {
	ChannelGroupID *uint
	AuthID         string
}

type ChannelGroupDetailRecord added in v1.0.44

type ChannelGroupDetailRecord struct {
	ID             uint           `` /* 162-byte string literal not displayed */
	ChannelGroupID uint           `` /* 160-byte string literal not displayed */
	AuthID         string         `gorm:"column:auth_id;not null;index:idx_channel_group_detail_auth_active_order,priority:1"`
	CreatedAt      time.Time      `gorm:"column:created_at"`
	UpdatedAt      time.Time      `gorm:"column:updated_at"`
	DeletedAt      gorm.DeletedAt `` /* 151-byte string literal not displayed */
}

func (ChannelGroupDetailRecord) TableName added in v1.0.44

func (ChannelGroupDetailRecord) TableName() string

TableName returns the database table name.

type ChannelGroupDetailUpdate added in v1.0.44

type ChannelGroupDetailUpdate struct {
	ChannelGroupID *uint
	AuthID         *string
}

type ChannelGroupRecord added in v1.0.44

type ChannelGroupRecord struct {
	ID          uint           `` /* 138-byte string literal not displayed */
	ChannelName string         `gorm:"column:channel_name;not null"`
	Disabled    bool           `gorm:"column:disabled;not null;default:false;index:idx_channel_group_enabled_order,priority:2"`
	CreatedAt   time.Time      `gorm:"column:created_at"`
	UpdatedAt   time.Time      `gorm:"column:updated_at"`
	DeletedAt   gorm.DeletedAt `` /* 127-byte string literal not displayed */
}

func (ChannelGroupRecord) TableName added in v1.0.44

func (ChannelGroupRecord) TableName() string

TableName returns the database table name.

type ChannelGroupUpdate added in v1.0.44

type ChannelGroupUpdate struct {
	ChannelName *string
	Disabled    *bool
}

type ClusterEventRecord

type ClusterEventRecord struct {
	ID         uint      `gorm:"column:id;primaryKey;autoIncrement"`
	Scope      string    `gorm:"column:scope"`
	Op         string    `gorm:"column:op"`
	EntityUUID string    `gorm:"column:entity_uuid"`
	Version    int64     `gorm:"column:version"`
	CreatedAt  time.Time `gorm:"column:created_at"`
}

func (ClusterEventRecord) TableName

func (ClusterEventRecord) TableName() string

TableName returns the database table name.

type ClusterNodeRecord

type ClusterNodeRecord struct {
	IP          string    `` /* 282-byte string literal not displayed */
	Port        int       `` /* 284-byte string literal not displayed */
	SecretHash  string    `gorm:"column:secret_hash;index:idx_cluster_auth_lookup,priority:2;index:idx_cluster_auth_live_first,priority:2"`
	IsMaster    bool      `gorm:"column:is_master;index:idx_cluster_master_nodes,priority:1;index:idx_cluster_master_current,priority:1"`
	ClientCount int       `gorm:"column:client_count;index:idx_cluster_live_schedule,priority:1"`
	StartedAt   time.Time `` /* 279-byte string literal not displayed */
	LastSeenAt  time.Time `` /* 291-byte string literal not displayed */
}

func (ClusterNodeRecord) TableName

func (ClusterNodeRecord) TableName() string

TableName returns the database table name.

type Config

type Config struct {
	PGSQL  PGSQLConfig  `yaml:"pgsql"`
	SQLite SQLiteConfig `yaml:"sqlite"`
	Node   NodeConfig   `yaml:"node"`
}

func LoadConfigOptional

func LoadConfigOptional(path string) (*Config, bool, error)

LoadConfigOptional loads a config optional.

func (*Config) DatabaseBackend added in v1.0.44

func (c *Config) DatabaseBackend() DatabaseBackend

DatabaseBackend returns the selected database backend.

func (*Config) Validate

func (c *Config) Validate() error

Validate validates validate.

type ConfigRecord

type ConfigRecord struct {
	Key       string    `gorm:"column:key;primaryKey"`
	Value     JSONB     `gorm:"column:value;not null"`
	Version   int64     `gorm:"column:version;not null;default:1"`
	CreatedAt time.Time `gorm:"column:created_at"`
	UpdatedAt time.Time `gorm:"column:updated_at"`
}

func (ConfigRecord) TableName

func (ConfigRecord) TableName() string

TableName returns the database table name.

type Coordinator

type Coordinator struct {
	// contains filtered or unexported fields
}

func NewCoordinator

func NewCoordinator(repo *Repository, node NodeIdentity, opts CoordinatorOptions) *Coordinator

NewCoordinator creates a new coordinator.

func (*Coordinator) CurrentMaster

func (c *Coordinator) CurrentMaster(ctx context.Context) (*ClusterNodeRecord, error)

CurrentMaster returns a current master.

func (*Coordinator) IsMaster

func (c *Coordinator) IsMaster() bool

IsMaster reports whether is master.

func (*Coordinator) NodeSecret

func (c *Coordinator) NodeSecret() string

NodeSecret handles a node secret.

func (*Coordinator) SetOnMasterChanged

func (c *Coordinator) SetOnMasterChanged(callback func(bool))

SetOnMasterChanged sets an on master changed.

func (*Coordinator) Start

func (c *Coordinator) Start(ctx context.Context) error

Start starts the process.

func (*Coordinator) UpdateClientCount

func (c *Coordinator) UpdateClientCount(ctx context.Context, clientCount int) error

UpdateClientCount stores the current active CPA client count for this node.

type CoordinatorOptions

type CoordinatorOptions struct {
	HeartbeatInterval time.Duration
	HeartbeatTimeout  time.Duration
	OnMasterChanged   func(bool)
}

type CredentialConfigCounts

type CredentialConfigCounts struct {
	GeminiKeys          int
	VertexKeys          int
	CodexKeys           int
	ClaudeKeys          int
	OpenAICompatibility int
}

CredentialConfigCounts reports how many auth-backed config entries were restored.

func ApplyCredentialConfigToRoot

func ApplyCredentialConfigToRoot(root map[string]any, auths []*coreauth.Auth) CredentialConfigCounts

ApplyCredentialConfigToRoot restores auth-backed config keys into a config root.

type DatabaseBackend added in v1.0.44

type DatabaseBackend string
const (
	DatabaseBackendPostgres DatabaseBackend = "postgres"
	DatabaseBackendSQLite   DatabaseBackend = "sqlite"
)

type EventWatcher

type EventWatcher struct {
	// contains filtered or unexported fields
}

func NewEventWatcher

func NewEventWatcher(repo *Repository, pollInterval time.Duration, onEvent func(context.Context, ClusterEventRecord) error) *EventWatcher

NewEventWatcher creates a new event watcher.

func NewEventWatcherFrom

func NewEventWatcherFrom(repo *Repository, pollInterval time.Duration, lastSeenID int64, onEvent func(context.Context, ClusterEventRecord) error) *EventWatcher

NewEventWatcherFrom creates a new event watcher from.

func (*EventWatcher) Start

func (w *EventWatcher) Start(ctx context.Context) error

Start starts the process.

type ExportOptions added in v1.0.44

type ExportOptions struct {
	OutputDir   string
	Repository  *Repository
	ConfigName  string
	AuthDirName string
}

type ExportStats added in v1.0.44

type ExportStats struct {
	ConfigBytes int
	AuthFiles   int
}

func ExportLocalState added in v1.0.44

func ExportLocalState(ctx context.Context, opts ExportOptions) (ExportStats, error)

ExportLocalState exports repository config and auth files into a local directory.

type ImportOptions added in v1.0.44

type ImportOptions struct {
	ConfigPath string
	AuthDir    string
	Repository *Repository
	Now        time.Time
}

type ImportStats added in v1.0.44

type ImportStats struct {
	Created     int
	Updated     int
	Unchanged   int
	Restored    int
	Overwritten int
	Skipped     int
}

func ImportLocalState added in v1.0.44

func ImportLocalState(ctx context.Context, opts ImportOptions) (ImportStats, error)

ImportLocalState imports config.yaml and auth-dir credentials into the repository.

type JSONB

type JSONB []byte

func NormalizeJSONB added in v1.0.44

func NormalizeJSONB(raw json.RawMessage) (*JSONB, error)

NormalizeJSONB validates and normalizes raw JSON for JSONB fields.

func (JSONB) GormDBDataType added in v1.0.44

func (JSONB) GormDBDataType(db *gorm.DB, _ *schema.Field) string

GormDBDataType returns the dialect-specific database type.

func (JSONB) GormDataType added in v1.0.44

func (JSONB) GormDataType() string

GormDataType returns the logical GORM data type.

func (JSONB) MarshalJSON

func (j JSONB) MarshalJSON() ([]byte, error)

MarshalJSON encodes a json.

func (*JSONB) Scan

func (j *JSONB) Scan(value any) error

Scan loads the value from database storage.

func (*JSONB) UnmarshalJSON

func (j *JSONB) UnmarshalJSON(data []byte) error

UnmarshalJSON decodes a json.

func (JSONB) Value

func (j JSONB) Value() (driver.Value, error)

Value converts the value for database storage.

type KVGetResult added in v1.0.44

type KVGetResult struct {
	Value []byte
	Found bool
}

type KVRecord added in v1.0.44

type KVRecord struct {
	Key       string     `gorm:"column:key;primaryKey"`
	Value     []byte     `gorm:"column:value;not null"`
	Version   int64      `gorm:"column:version;not null;default:1"`
	ExpiresAt *time.Time `gorm:"column:expires_at;index"`
	CreatedAt time.Time  `gorm:"column:created_at"`
	UpdatedAt time.Time  `gorm:"column:updated_at"`
}

func (KVRecord) TableName added in v1.0.44

func (KVRecord) TableName() string

TableName returns the database table name.

type KVSetMode added in v1.0.44

type KVSetMode string
const (
	KVSetModeAlways KVSetMode = ""
	KVSetModeNX     KVSetMode = "nx"
	KVSetModeXX     KVSetMode = "xx"
)

type ModelGroupDetailFilter added in v1.0.44

type ModelGroupDetailFilter struct {
	ModelGroupID *uint
	ModelID      string
}

type ModelGroupDetailRecord added in v1.0.44

type ModelGroupDetailRecord struct {
	ID           uint           `` /* 212-byte string literal not displayed */
	ModelGroupID uint           `` /* 208-byte string literal not displayed */
	ModelID      string         `gorm:"column:model_id;not null;index:idx_model_group_detail_model_active_order,priority:1"`
	CreatedAt    time.Time      `gorm:"column:created_at"`
	UpdatedAt    time.Time      `gorm:"column:updated_at"`
	DeletedAt    gorm.DeletedAt `` /* 201-byte string literal not displayed */
}

func (ModelGroupDetailRecord) TableName added in v1.0.44

func (ModelGroupDetailRecord) TableName() string

TableName returns the database table name.

type ModelGroupDetailUpdate added in v1.0.44

type ModelGroupDetailUpdate struct {
	ModelGroupID *uint
	ModelID      *string
}

type ModelGroupRecord added in v1.0.44

type ModelGroupRecord struct {
	ID        uint           `` /* 134-byte string literal not displayed */
	GroupName string         `gorm:"column:group_name;not null"`
	Disabled  bool           `gorm:"column:disabled;not null;default:false;index:idx_model_group_enabled_order,priority:2"`
	CreatedAt time.Time      `gorm:"column:created_at"`
	UpdatedAt time.Time      `gorm:"column:updated_at"`
	DeletedAt gorm.DeletedAt `gorm:"column:deleted_at;index;index:idx_model_group_active_order,priority:1;index:idx_model_group_enabled_order,priority:1"`
}

func (ModelGroupRecord) TableName added in v1.0.44

func (ModelGroupRecord) TableName() string

TableName returns the database table name.

type ModelGroupUpdate added in v1.0.44

type ModelGroupUpdate struct {
	GroupName *string
	Disabled  *bool
}

type NodeConfig

type NodeConfig struct {
	ExternalIP        string        `yaml:"external-ip"`
	ExternalPort      int           `yaml:"external-port"`
	Port              int           `yaml:"port"`
	HeartbeatInterval time.Duration `yaml:"heartbeat-interval"`
	HeartbeatTimeout  time.Duration `yaml:"heartbeat-timeout"`
	EventPollInterval time.Duration `yaml:"event-poll-interval"`
}

type NodeIdentity

type NodeIdentity struct {
	IP        string
	Port      int
	Secret    string
	StartedAt time.Time
}

type OAuthSessionRecord

type OAuthSessionRecord struct {
	State       string     `gorm:"column:state;primaryKey"`
	Provider    string     `gorm:"column:provider;index"`
	Status      string     `gorm:"column:status"`
	Error       string     `gorm:"column:error"`
	Data        JSONB      `gorm:"column:data"`
	CreatedAt   time.Time  `gorm:"column:created_at"`
	UpdatedAt   time.Time  `gorm:"column:updated_at"`
	ExpiresAt   time.Time  `gorm:"column:expires_at;index"`
	CompletedAt *time.Time `gorm:"column:completed_at"`
}

func NewOAuthSessionRecord

func NewOAuthSessionRecord(provider, state string, data map[string]any, now time.Time) (*OAuthSessionRecord, error)

NewOAuthSessionRecord creates a new o auth session record.

func (OAuthSessionRecord) TableName

func (OAuthSessionRecord) TableName() string

TableName returns the database table name.

type PGSQLConfig

type PGSQLConfig struct {
	Host     string `yaml:"host"`
	Port     int    `yaml:"port"`
	User     string `yaml:"user"`
	Password string `yaml:"password"`
	Passowrd string `yaml:"passowrd"`
	Database string `yaml:"database"`
	SSLMode  string `yaml:"sslmode"`
}

func (PGSQLConfig) Configured added in v1.0.44

func (c PGSQLConfig) Configured() bool

Configured reports whether PostgreSQL has cluster database settings.

func (PGSQLConfig) DSN

func (c PGSQLConfig) DSN() (string, error)

DSN returns the PostgreSQL connection string.

func (PGSQLConfig) Validate

func (c PGSQLConfig) Validate() error

Validate validates validate.

type PluginStatusRecord added in v1.0.44

type PluginStatusRecord struct {
	NodeType      string     `gorm:"column:node_type;primaryKey;size:16;index:idx_plugin_status_node,priority:1;index:idx_plugin_status_plugin,priority:1"`
	NodeID        string     `gorm:"column:node_id;primaryKey;size:128;index:idx_plugin_status_node,priority:2"`
	PluginID      string     `gorm:"column:plugin_id;primaryKey;size:128;index:idx_plugin_status_plugin,priority:2"`
	TaskID        uint       `gorm:"column:task_id;index"`
	ClientIP      string     `gorm:"column:client_ip;index"`
	SchemaVersion int        `gorm:"column:schema_version"`
	Task          string     `gorm:"column:task"`
	TaskStatus    string     `gorm:"column:task_status;index"`
	Phase         string     `gorm:"column:phase;index"`
	OK            bool       `gorm:"column:ok;index"`
	TaskError     string     `gorm:"column:task_error;type:text"`
	StartedAt     time.Time  `gorm:"column:started_at"`
	FinishedAt    *time.Time `gorm:"column:finished_at"`
	ReportedAt    time.Time  `gorm:"column:reported_at;index"`
	GOOS          string     `gorm:"column:goos;index"`
	GOARCH        string     `gorm:"column:goarch;index"`
	Variant       string     `gorm:"column:variant"`
	Version       string     `gorm:"column:version"`
	ReleaseTag    string     `gorm:"column:release_tag"`
	Repository    string     `gorm:"column:repository;type:text"`
	InstallStatus string     `gorm:"column:install_status;index"`
	LoadStatus    string     `gorm:"column:load_status;index"`
	Path          string     `gorm:"column:path;type:text"`
	Skipped       bool       `gorm:"column:skipped"`
	Overwritten   bool       `gorm:"column:overwritten"`
	PluginError   string     `gorm:"column:plugin_error;type:text"`
	CreatedAt     time.Time  `gorm:"column:created_at"`
	UpdatedAt     time.Time  `gorm:"column:updated_at"`
}

func (PluginStatusRecord) TableName added in v1.0.44

func (PluginStatusRecord) TableName() string

TableName returns the database table name.

type PluginTaskRecord added in v1.0.44

type PluginTaskRecord struct {
	ID             uint      `gorm:"column:id;primaryKey;autoIncrement"`
	Operation      string    `gorm:"column:operation;size:32;index:idx_plugin_tasks_operation"`
	PluginID       string    `gorm:"column:plugin_id;size:128;index:idx_plugin_tasks_plugin"`
	TargetNodeType string    `gorm:"column:target_node_type;size:16;index:idx_plugin_tasks_target,priority:1"`
	TargetNodeID   string    `gorm:"column:target_node_id;size:128;index:idx_plugin_tasks_target,priority:2"`
	CreatedAt      time.Time `gorm:"column:created_at"`
	UpdatedAt      time.Time `gorm:"column:updated_at"`
}

func (PluginTaskRecord) TableName added in v1.0.44

func (PluginTaskRecord) TableName() string

TableName returns the database table name.

type ProxyPoolPatch added in v1.0.44

type ProxyPoolPatch struct {
	Name     *string
	ProxyURL *string
	Enabled  *bool
	Scope    *string
	Priority *int
	Note     *string
}

type ProxyPoolRecord added in v1.0.44

type ProxyPoolRecord struct {
	ID             string         `gorm:"column:id;primaryKey"`
	Name           string         `gorm:"column:name;not null"`
	ProxyURL       string         `gorm:"column:proxy_url;not null"`
	Enabled        bool           `gorm:"column:enabled;not null;default:true"`
	Scope          string         `gorm:"column:scope;not null;default:global;index:idx_proxy_pool_scope_priority,priority:1"`
	Priority       int            `gorm:"column:priority;not null;default:0;index:idx_proxy_pool_scope_priority,priority:2"`
	LastTestedAt   *time.Time     `gorm:"column:last_tested_at"`
	LastTestResult string         `gorm:"column:last_test_result;not null;default:untested"`
	Note           string         `gorm:"column:note;type:text"`
	CreatedAt      time.Time      `gorm:"column:created_at;index:idx_proxy_pool_scope_priority,priority:3"`
	UpdatedAt      time.Time      `gorm:"column:updated_at"`
	DeletedAt      gorm.DeletedAt `gorm:"column:deleted_at;index"`
}

func (ProxyPoolRecord) TableName added in v1.0.44

func (ProxyPoolRecord) TableName() string

type ProxyPoolUpdate added in v1.0.44

type ProxyPoolUpdate struct {
	Name     string
	ProxyURL string
	Enabled  *bool
	Scope    string
	Priority int
	Note     string
}

type RESPHandler

type RESPHandler struct {
	// contains filtered or unexported fields
}

func NewRESPHandler

func NewRESPHandler(coordinator *Coordinator, refresh *RefreshController, repo *Repository) *RESPHandler

NewRESPHandler creates a new resp handler.

func (*RESPHandler) Handle

func (h *RESPHandler) Handle(ctx context.Context, args []string, remoteIP string) ([]byte, error)

Handle handles handle.

func (*RESPHandler) RequestClientCertificate added in v1.0.44

func (h *RESPHandler) RequestClientCertificate(ctx context.Context, certificateID string, enrollmentSecret string, csrPEM []byte) ([]byte, error)

RequestClientCertificate signs a pending client certificate request.

func (*RESPHandler) UpdateClientCount

func (h *RESPHandler) UpdateClientCount(ctx context.Context, clientCount int) error

UpdateClientCount stores the current active CPA client count for this node.

type RefreshController

type RefreshController struct {
	// contains filtered or unexported fields
}

func NewRefreshController

func NewRefreshController(coordinator *Coordinator, runtime *home.Runtime, repo *Repository, forwardTLSConfig *tls.Config) *RefreshController

NewRefreshController creates a new refresh controller.

func (*RefreshController) CanAutoRefresh

func (c *RefreshController) CanAutoRefresh() bool

CanAutoRefresh reports whether can auto refresh.

func (*RefreshController) OnMasterChanged

func (c *RefreshController) OnMasterChanged(isMaster bool)

OnMasterChanged handles an on master changed.

func (*RefreshController) RefreshNow

func (c *RefreshController) RefreshNow(ctx context.Context, authIndex string) ([]byte, error)

RefreshNow refreshes refresh now.

type Repository

type Repository struct {
	// contains filtered or unexported fields
}

func NewRepository

func NewRepository(db *gorm.DB) *Repository

NewRepository creates a new repository.

func (*Repository) AllowedAuthIDsForAPIKey added in v1.0.44

func (r *Repository) AllowedAuthIDsForAPIKey(ctx context.Context, apiKey string) ([]string, error)

AllowedAuthIDsForAPIKey returns auth IDs allowed by the API key channel bindings.

func (*Repository) AllowedDispatchIDsForAPIKey added in v1.0.44

func (r *Repository) AllowedDispatchIDsForAPIKey(ctx context.Context, apiKey string) ([]string, []string, error)

AllowedDispatchIDsForAPIKey returns auth and model IDs allowed by API-key bindings.

func (*Repository) AllowedModelIDsForAPIKey added in v1.0.44

func (r *Repository) AllowedModelIDsForAPIKey(ctx context.Context, apiKey string) ([]string, error)

AllowedModelIDsForAPIKey returns model IDs allowed by the API key model group bindings.

func (*Repository) AppendAppLog added in v1.0.44

func (r *Repository) AppendAppLog(ctx context.Context, clientIP string, homeIP string, payload string) (*AppLogRecord, error)

AppendAppLog appends a CPA application log record.

func (*Repository) AppendEvent

func (r *Repository) AppendEvent(ctx context.Context, scope, op, entity string, version int64) error

AppendEvent appends an event.

func (*Repository) AppendUsage

func (r *Repository) AppendUsage(ctx context.Context, payload string, homeIP string) (*UsageRecord, error)

AppendUsage appends an usage.

func (*Repository) ApplyBillingBalanceRecord added in v1.0.44

func (r *Repository) ApplyBillingBalanceRecord(ctx context.Context, update BillingBalanceUpdate) (*BillingBalanceRecord, error)

func (*Repository) BillingOverview added in v1.0.44

func (r *Repository) BillingOverview(ctx context.Context, query BillingOverviewQuery) (BillingOverview, error)

func (*Repository) ClusterCAKeyPair added in v1.0.44

func (r *Repository) ClusterCAKeyPair(ctx context.Context) (*x509.Certificate, *rsa.PrivateKey, error)

ClusterCAKeyPair returns the cluster root CA certificate and private key.

func (*Repository) CompleteOAuthSession

func (r *Repository) CompleteOAuthSession(ctx context.Context, state string) error

CompleteOAuthSession handles a complete o auth session.

func (*Repository) CreateAPIKeyForUser added in v1.0.44

func (r *Repository) CreateAPIKeyForUser(ctx context.Context, userID uint, update APIKeyUserUpdate) (*APIKeyRecord, error)

CreateAPIKeyForUser creates an API key owned by one user.

func (*Repository) CreateBillingModelPrice added in v1.0.44

func (r *Repository) CreateBillingModelPrice(ctx context.Context, update BillingModelPriceUpdate) (*BillingModelPriceRecord, error)

func (*Repository) CreateChannelGroup added in v1.0.44

func (r *Repository) CreateChannelGroup(ctx context.Context, channelName string, disabled bool) (*ChannelGroupRecord, error)

CreateChannelGroup creates a channel group.

func (*Repository) CreateChannelGroupDetail added in v1.0.44

func (r *Repository) CreateChannelGroupDetail(ctx context.Context, channelGroupID uint, authID string) (*ChannelGroupDetailRecord, error)

CreateChannelGroupDetail creates a channel group detail.

func (*Repository) CreateHomeJWT added in v1.0.44

func (r *Repository) CreateHomeJWT(ctx context.Context, certificateID string, ip string, port int, enrollmentSecret string) (string, error)

CreateHomeJWT creates a signed Home JWT with the client certificate id and target node.

func (*Repository) CreateModelGroup added in v1.0.44

func (r *Repository) CreateModelGroup(ctx context.Context, groupName string, disabled bool) (*ModelGroupRecord, error)

CreateModelGroup creates a model group.

func (*Repository) CreateModelGroupDetail added in v1.0.44

func (r *Repository) CreateModelGroupDetail(ctx context.Context, modelGroupID uint, modelID string) (*ModelGroupDetailRecord, error)

CreateModelGroupDetail creates a model group detail.

func (*Repository) CreatePendingClientCertificate added in v1.0.44

func (r *Repository) CreatePendingClientCertificate(ctx context.Context) (string, string, error)

CreatePendingClientCertificate creates an empty client certificate slot.

func (*Repository) CreatePluginTask added in v1.0.44

func (r *Repository) CreatePluginTask(ctx context.Context, task node.PluginTask) (node.PluginTask, error)

CreatePluginTask stores a plugin task and emits a cluster event to wake peers.

func (*Repository) CreateProxyPoolItem added in v1.0.44

func (r *Repository) CreateProxyPoolItem(ctx context.Context, update ProxyPoolUpdate) (*ProxyPoolRecord, error)

func (*Repository) CreateUser added in v1.0.44

func (r *Repository) CreateUser(ctx context.Context, update UserUpdate) (*UserRecord, error)

CreateUser creates a user.

func (*Repository) CurrentMasterNode added in v1.0.44

func (r *Repository) CurrentMasterNode(ctx context.Context) (*ClusterNodeRecord, error)

CurrentMasterNode returns the current master node if one is recorded.

func (*Repository) DeleteAPIKeyForUser added in v1.0.44

func (r *Repository) DeleteAPIKeyForUser(ctx context.Context, userID uint, id uint, apiKey string) error

DeleteAPIKeyForUser deletes one API key owned by one user.

func (*Repository) DeleteBillingModelPrice added in v1.0.44

func (r *Repository) DeleteBillingModelPrice(ctx context.Context, id string) error

func (*Repository) DeleteChannelGroup added in v1.0.44

func (r *Repository) DeleteChannelGroup(ctx context.Context, id uint) error

DeleteChannelGroup deletes a channel group and its details.

func (*Repository) DeleteChannelGroupDetail added in v1.0.44

func (r *Repository) DeleteChannelGroupDetail(ctx context.Context, id uint) error

DeleteChannelGroupDetail deletes a channel group detail.

func (*Repository) DeleteModelGroup added in v1.0.44

func (r *Repository) DeleteModelGroup(ctx context.Context, id uint) error

DeleteModelGroup deletes a model group and its details.

func (*Repository) DeleteModelGroupDetail added in v1.0.44

func (r *Repository) DeleteModelGroupDetail(ctx context.Context, id uint) error

DeleteModelGroupDetail deletes a model group detail.

func (*Repository) DeleteProxyPoolItem added in v1.0.44

func (r *Repository) DeleteProxyPoolItem(ctx context.Context, id string) error

func (*Repository) DeleteUser added in v1.0.44

func (r *Repository) DeleteUser(ctx context.Context, id uint) error

DeleteUser deletes a user.

func (*Repository) EnsureClusterCertificates added in v1.0.44

func (r *Repository) EnsureClusterCertificates(ctx context.Context, ip string) (*tls.Config, error)

EnsureClusterCertificates makes sure the cluster CA and this node server certificate exist.

func (*Repository) GetAuth

func (r *Repository) GetAuth(ctx context.Context, uuid string) (*coreauth.Auth, *AuthRecord, error)

GetAuth returns an auth.

func (*Repository) GetBillingModelPrice added in v1.0.44

func (r *Repository) GetBillingModelPrice(ctx context.Context, id string) (*BillingModelPriceRecord, error)

func (*Repository) GetChannelGroup added in v1.0.44

func (r *Repository) GetChannelGroup(ctx context.Context, id uint) (*ChannelGroupRecord, error)

GetChannelGroup returns a channel group by ID.

func (*Repository) GetChannelGroupDetail added in v1.0.44

func (r *Repository) GetChannelGroupDetail(ctx context.Context, id uint) (*ChannelGroupDetailRecord, error)

GetChannelGroupDetail returns a channel group detail by ID.

func (*Repository) GetModelGroup added in v1.0.44

func (r *Repository) GetModelGroup(ctx context.Context, id uint) (*ModelGroupRecord, error)

GetModelGroup returns a model group by ID.

func (*Repository) GetModelGroupDetail added in v1.0.44

func (r *Repository) GetModelGroupDetail(ctx context.Context, id uint) (*ModelGroupDetailRecord, error)

GetModelGroupDetail returns a model group detail by ID.

func (*Repository) GetOAuthSession

func (r *Repository) GetOAuthSession(ctx context.Context, state string) (*OAuthSessionRecord, error)

GetOAuthSession returns an o auth session.

func (*Repository) GetProxyPoolItem added in v1.0.44

func (r *Repository) GetProxyPoolItem(ctx context.Context, id string) (*ProxyPoolRecord, error)

func (*Repository) GetUser added in v1.0.44

func (r *Repository) GetUser(ctx context.Context, id uint) (*UserRecord, error)

GetUser returns a user by ID.

func (*Repository) GetUserByUsername added in v1.0.44

func (r *Repository) GetUserByUsername(ctx context.Context, username string) (*UserRecord, error)

GetUserByUsername returns a user by username.

func (*Repository) KVDel added in v1.0.44

func (r *Repository) KVDel(ctx context.Context, keys []string) (int64, error)

KVDel deletes active keys and returns the active deletion count.

func (*Repository) KVExpire added in v1.0.44

func (r *Repository) KVExpire(ctx context.Context, key string, ttl time.Duration) (bool, error)

KVExpire updates the TTL for an active key.

func (*Repository) KVGet added in v1.0.44

func (r *Repository) KVGet(ctx context.Context, key string) ([]byte, bool, error)

KVGet returns the active value for a key.

func (*Repository) KVIncrBy added in v1.0.44

func (r *Repository) KVIncrBy(ctx context.Context, key string, delta int64) (int64, error)

KVIncrBy increments a decimal integer value.

func (*Repository) KVMGet added in v1.0.44

func (r *Repository) KVMGet(ctx context.Context, keys []string) ([]KVGetResult, error)

KVMGet returns values in the same order as the requested keys.

func (*Repository) KVMSet added in v1.0.44

func (r *Repository) KVMSet(ctx context.Context, pairs map[string][]byte) error

KVMSet atomically writes key/value pairs without TTL.

func (*Repository) KVPurgeExpired added in v1.0.44

func (r *Repository) KVPurgeExpired(ctx context.Context, now time.Time, limit int) (int64, error)

KVPurgeExpired deletes expired rows up to limit.

func (*Repository) KVSet added in v1.0.44

func (r *Repository) KVSet(ctx context.Context, key string, value []byte, ttl time.Duration, mode KVSetMode) (bool, error)

KVSet writes a value using the requested conditional mode.

func (*Repository) KVTTL added in v1.0.44

func (r *Repository) KVTTL(ctx context.Context, key string) (int64, error)

KVTTL returns Redis-compatible TTL seconds.

func (*Repository) ListAPIKeyEntries added in v1.0.44

func (r *Repository) ListAPIKeyEntries(ctx context.Context) ([]APIKeyEntry, error)

ListAPIKeyEntries returns API key rows with group bindings.

func (*Repository) ListAPIKeyRecordsForUser added in v1.0.44

func (r *Repository) ListAPIKeyRecordsForUser(ctx context.Context, userID uint) ([]APIKeyRecord, error)

ListAPIKeyRecordsForUser returns active API key rows owned by one user.

func (*Repository) ListAppLogs added in v1.0.44

func (r *Repository) ListAppLogs(ctx context.Context, opts AppLogQuery) (AppLogQueryResult, error)

ListAppLogs returns application log records from the database.

func (*Repository) ListAuthIndex

func (r *Repository) ListAuthIndex(ctx context.Context) ([]AuthIndex, error)

ListAuthIndex returns an auth index.

func (*Repository) ListAuths

func (r *Repository) ListAuths(ctx context.Context) ([]*coreauth.Auth, error)

ListAuths returns an auths.

func (*Repository) ListBillingBalanceRecords added in v1.0.44

func (r *Repository) ListBillingBalanceRecords(ctx context.Context, query BillingBalanceQuery) (BillingBalanceResult, error)

func (*Repository) ListBillingCharges added in v1.0.44

func (r *Repository) ListBillingCharges(ctx context.Context, query BillingChargeQuery) (BillingChargeResult, error)

func (*Repository) ListBillingModelPrices added in v1.0.44

func (r *Repository) ListBillingModelPrices(ctx context.Context, query BillingModelPriceQuery) ([]BillingModelPriceRecord, error)

func (*Repository) ListChannelGroupDetails added in v1.0.44

func (r *Repository) ListChannelGroupDetails(ctx context.Context, filter ChannelGroupDetailFilter) ([]ChannelGroupDetailRecord, error)

ListChannelGroupDetails returns channel group details.

func (*Repository) ListChannelGroups added in v1.0.44

func (r *Repository) ListChannelGroups(ctx context.Context) ([]ChannelGroupRecord, error)

ListChannelGroups returns channel groups.

func (*Repository) ListLiveClusterNodes

func (r *Repository) ListLiveClusterNodes(ctx context.Context, cutoff time.Time) ([]ClusterNodeRecord, error)

ListLiveClusterNodes returns live cluster nodes.

func (*Repository) ListModelGroupDetails added in v1.0.44

func (r *Repository) ListModelGroupDetails(ctx context.Context, filter ModelGroupDetailFilter) ([]ModelGroupDetailRecord, error)

ListModelGroupDetails returns model group details.

func (*Repository) ListModelGroups added in v1.0.44

func (r *Repository) ListModelGroups(ctx context.Context) ([]ModelGroupRecord, error)

ListModelGroups returns model groups.

func (*Repository) ListPendingPluginTasks added in v1.0.44

func (r *Repository) ListPendingPluginTasks(ctx context.Context, nodeType string, nodeID string) ([]node.PluginTask, error)

ListPendingPluginTasks returns tasks that the given node has not acked yet.

func (*Repository) ListPluginStatuses added in v1.0.44

func (r *Repository) ListPluginStatuses(ctx context.Context, nodeType string) ([]node.PluginTaskStatus, error)

ListPluginStatuses returns the latest plugin task statuses grouped by node and report.

func (*Repository) ListProxyPoolItems added in v1.0.44

func (r *Repository) ListProxyPoolItems(ctx context.Context) ([]ProxyPoolRecord, error)

func (*Repository) ListUsers added in v1.0.44

func (r *Repository) ListUsers(ctx context.Context) ([]UserRecord, error)

ListUsers returns users.

func (*Repository) LiveClusterNodeByIPAndSecret

func (r *Repository) LiveClusterNodeByIPAndSecret(ctx context.Context, ip string, secret string, cutoff time.Time) (*ClusterNodeRecord, error)

LiveClusterNodeByIPAndSecret handles a live cluster node by ip and secret.

func (*Repository) LoadConfigAsRuntimeConfig

func (r *Repository) LoadConfigAsRuntimeConfig(ctx context.Context) (*appconfig.Config, []byte, error)

LoadConfigAsRuntimeConfig loads a config as runtime config.

func (*Repository) LoadConfigSnapshot

func (r *Repository) LoadConfigSnapshot(ctx context.Context) (map[string]json.RawMessage, error)

LoadConfigSnapshot loads a config snapshot.

func (*Repository) MarkProxyPoolTestResult added in v1.0.44

func (r *Repository) MarkProxyPoolTestResult(ctx context.Context, id string, result string, testedAt time.Time) (*ProxyPoolRecord, error)

func (*Repository) MaxEventID

func (r *Repository) MaxEventID(ctx context.Context) (int64, error)

MaxEventID handles a max event id.

func (*Repository) MergeOAuthSessionData

func (r *Repository) MergeOAuthSessionData(ctx context.Context, state string, values map[string]any) error

MergeOAuthSessionData merges an o auth session data.

func (*Repository) PatchProxyPoolItem added in v1.0.44

func (r *Repository) PatchProxyPoolItem(ctx context.Context, id string, patch ProxyPoolPatch) (*ProxyPoolRecord, error)

func (*Repository) ReplaceAPIKeyEntries added in v1.0.44

func (r *Repository) ReplaceAPIKeyEntries(ctx context.Context, entries []APIKeyEntryUpdate) (APIKeyUpsertStats, error)

ReplaceAPIKeyEntries replaces active API key rows and updates explicit channel bindings.

func (*Repository) ReplaceConfigSnapshot

func (r *Repository) ReplaceConfigSnapshot(ctx context.Context, values map[string]any) error

ReplaceConfigSnapshot handles a replace config snapshot.

func (*Repository) ReplaceConfigSnapshotAndCreatePluginTask added in v1.0.44

func (r *Repository) ReplaceConfigSnapshotAndCreatePluginTask(ctx context.Context, values map[string]any, task node.PluginTask) (node.PluginTask, error)

ReplaceConfigSnapshotAndCreatePluginTask replaces config and creates a plugin task atomically.

func (*Repository) ReplacePluginStatus added in v1.0.44

func (r *Repository) ReplacePluginStatus(ctx context.Context, nodeType string, status node.PluginTaskStatus) error

ReplacePluginStatus stores the latest plugin task status for one node.

func (*Repository) SetOAuthSessionError

func (r *Repository) SetOAuthSessionError(ctx context.Context, state string, message string) error

SetOAuthSessionError sets an o auth session error.

func (*Repository) SignClientCertificateRequest added in v1.0.44

func (r *Repository) SignClientCertificateRequest(ctx context.Context, certificateID string, enrollmentSecret string, csrPEM []byte) ([]byte, []byte, error)

SignClientCertificateRequest signs a client CSR exactly once and returns the client certificate plus CA.

func (*Repository) SignClientCertificateRequestJSON added in v1.0.44

func (r *Repository) SignClientCertificateRequestJSON(ctx context.Context, certificateID string, enrollmentSecret string, csrPEM []byte) ([]byte, error)

SignClientCertificateRequestJSON signs a CSR and returns the RESP JSON payload.

func (*Repository) SoftDeleteAuth

func (r *Repository) SoftDeleteAuth(ctx context.Context, uuid string) error

SoftDeleteAuth handles a soft delete auth.

func (*Repository) UpdateAPIKeyBindings added in v1.0.44

func (r *Repository) UpdateAPIKeyBindings(ctx context.Context, apiKey string, userID *uint, channels *[]uint, modelGroups *[]uint) (*APIKeyRecord, error)

UpdateAPIKeyBindings updates user and group bindings for one API key.

func (*Repository) UpdateAPIKeyChannels added in v1.0.44

func (r *Repository) UpdateAPIKeyChannels(ctx context.Context, apiKey string, channels []uint) (*APIKeyRecord, error)

UpdateAPIKeyChannels updates channel bindings for one API key.

func (*Repository) UpdateAPIKeyForUser added in v1.0.44

func (r *Repository) UpdateAPIKeyForUser(ctx context.Context, userID uint, id uint, apiKey string, update APIKeyUserUpdate) (*APIKeyRecord, error)

UpdateAPIKeyForUser updates one API key owned by one user.

func (*Repository) UpdateBillingModelPrice added in v1.0.44

func (r *Repository) UpdateBillingModelPrice(ctx context.Context, id string, patch BillingModelPricePatch) (*BillingModelPriceRecord, error)

func (*Repository) UpdateChannelGroup added in v1.0.44

func (r *Repository) UpdateChannelGroup(ctx context.Context, id uint, update ChannelGroupUpdate) (*ChannelGroupRecord, error)

UpdateChannelGroup updates a channel group.

func (*Repository) UpdateChannelGroupDetail added in v1.0.44

func (r *Repository) UpdateChannelGroupDetail(ctx context.Context, id uint, update ChannelGroupDetailUpdate) (*ChannelGroupDetailRecord, error)

UpdateChannelGroupDetail updates a channel group detail.

func (*Repository) UpdateModelGroup added in v1.0.44

func (r *Repository) UpdateModelGroup(ctx context.Context, id uint, update ModelGroupUpdate) (*ModelGroupRecord, error)

UpdateModelGroup updates a model group.

func (*Repository) UpdateModelGroupDetail added in v1.0.44

func (r *Repository) UpdateModelGroupDetail(ctx context.Context, id uint, update ModelGroupDetailUpdate) (*ModelGroupDetailRecord, error)

UpdateModelGroupDetail updates a model group detail.

func (*Repository) UpdateProxyPoolItem added in v1.0.44

func (r *Repository) UpdateProxyPoolItem(ctx context.Context, id string, update ProxyPoolUpdate) (*ProxyPoolRecord, error)

func (*Repository) UpdateUser added in v1.0.44

func (r *Repository) UpdateUser(ctx context.Context, id uint, update UserUpdate) (*UserRecord, error)

UpdateUser updates a user.

func (*Repository) UpsertAPIKeys added in v1.0.44

func (r *Repository) UpsertAPIKeys(ctx context.Context, keys []string) (APIKeyUpsertStats, error)

UpsertAPIKeys inserts or restores API keys without deleting keys missing from the input.

func (*Repository) UpsertAuth

func (r *Repository) UpsertAuth(ctx context.Context, auth *coreauth.Auth, op string) (*AuthRecord, error)

UpsertAuth inserts or updates an auth.

func (*Repository) UpsertAuthWithResult added in v1.0.44

func (r *Repository) UpsertAuthWithResult(ctx context.Context, auth *coreauth.Auth, op string) (*AuthRecord, UpsertResult, error)

UpsertAuthWithResult inserts or updates an auth and reports the mutation result.

func (*Repository) UpsertConfigValue

func (r *Repository) UpsertConfigValue(ctx context.Context, key string, value any) error

UpsertConfigValue inserts or updates a config value.

func (*Repository) UpsertConfigValueWithResult added in v1.0.44

func (r *Repository) UpsertConfigValueWithResult(ctx context.Context, key string, value any) (UpsertResult, error)

UpsertConfigValueWithResult inserts or updates a config value and reports the mutation result.

func (*Repository) UpsertOAuthSession

func (r *Repository) UpsertOAuthSession(ctx context.Context, record *OAuthSessionRecord) error

UpsertOAuthSession inserts or updates an o auth session.

func (*Repository) ValidateAPIKey added in v1.0.44

func (r *Repository) ValidateAPIKey(ctx context.Context, apiKey string) (bool, error)

ValidateAPIKey reports whether an API key exists as an active record.

func (*Repository) WithAuthRefreshLock

func (r *Repository) WithAuthRefreshLock(ctx context.Context, uuid string, fn func(tx *Repository, auth *coreauth.Auth) (*coreauth.Auth, error)) (*coreauth.Auth, error)

WithAuthRefreshLock applies the auth refresh lock option.

type RuntimeAdapter

type RuntimeAdapter struct {
	// contains filtered or unexported fields
}

func NewRuntimeAdapter

func NewRuntimeAdapter(repo *Repository, homeIP string) *RuntimeAdapter

NewRuntimeAdapter creates a new runtime adapter.

func (*RuntimeAdapter) AllowedAuthIDsForAPIKey added in v1.0.44

func (a *RuntimeAdapter) AllowedAuthIDsForAPIKey(ctx context.Context, apiKey string) ([]string, error)

AllowedAuthIDsForAPIKey returns auth IDs allowed by API-key channel bindings.

func (*RuntimeAdapter) AllowedDispatchIDsForAPIKey added in v1.0.44

func (a *RuntimeAdapter) AllowedDispatchIDsForAPIKey(ctx context.Context, apiKey string) ([]string, []string, error)

AllowedDispatchIDsForAPIKey returns auth and model IDs allowed by API-key bindings.

func (*RuntimeAdapter) AllowedModelIDsForAPIKey added in v1.0.44

func (a *RuntimeAdapter) AllowedModelIDsForAPIKey(ctx context.Context, apiKey string) ([]string, error)

AllowedModelIDsForAPIKey returns model IDs allowed by API-key model group bindings.

func (*RuntimeAdapter) ApplyEvent

func (a *RuntimeAdapter) ApplyEvent(ctx context.Context, event ClusterEventRecord) error

ApplyEvent updates apply event.

func (*RuntimeAdapter) Delete

func (a *RuntimeAdapter) Delete(ctx context.Context, id string) error

Delete handles delete.

func (*RuntimeAdapter) Enabled

func (a *RuntimeAdapter) Enabled() bool

Enabled handles an enabled.

func (*RuntimeAdapter) GetFullAuth

func (a *RuntimeAdapter) GetFullAuth(ctx context.Context, uuid string) (*coreauth.Auth, error)

GetFullAuth returns a full auth.

func (*RuntimeAdapter) InvalidateFullAuth

func (a *RuntimeAdapter) InvalidateFullAuth(uuid string)

InvalidateFullAuth invalidates a full auth.

func (*RuntimeAdapter) KVDel added in v1.0.44

func (a *RuntimeAdapter) KVDel(ctx context.Context, keys []string) (int64, error)

KVDel deletes active KV values.

func (*RuntimeAdapter) KVExpire added in v1.0.44

func (a *RuntimeAdapter) KVExpire(ctx context.Context, key string, ttl time.Duration) (bool, error)

KVExpire updates a KV TTL.

func (*RuntimeAdapter) KVGet added in v1.0.44

func (a *RuntimeAdapter) KVGet(ctx context.Context, key string) ([]byte, bool, error)

KVGet returns an active KV value.

func (*RuntimeAdapter) KVIncrBy added in v1.0.44

func (a *RuntimeAdapter) KVIncrBy(ctx context.Context, key string, delta int64) (int64, error)

KVIncrBy increments a KV integer.

func (*RuntimeAdapter) KVMGet added in v1.0.44

func (a *RuntimeAdapter) KVMGet(ctx context.Context, keys []string) ([]home.KVGetResult, error)

KVMGet returns KV values in request order.

func (*RuntimeAdapter) KVMSet added in v1.0.44

func (a *RuntimeAdapter) KVMSet(ctx context.Context, pairs map[string][]byte) error

KVMSet atomically writes KV values.

func (*RuntimeAdapter) KVPurgeExpired added in v1.0.44

func (a *RuntimeAdapter) KVPurgeExpired(ctx context.Context, now time.Time, limit int) (int64, error)

KVPurgeExpired deletes expired KV rows.

func (*RuntimeAdapter) KVSet added in v1.0.44

func (a *RuntimeAdapter) KVSet(ctx context.Context, key string, value []byte, ttl time.Duration, mode string) (bool, error)

KVSet writes a KV value.

func (*RuntimeAdapter) KVTTL added in v1.0.44

func (a *RuntimeAdapter) KVTTL(ctx context.Context, key string) (int64, error)

KVTTL returns a Redis-compatible KV TTL.

func (*RuntimeAdapter) List

func (a *RuntimeAdapter) List(ctx context.Context) ([]*coreauth.Auth, error)

List returns the available entries.

func (*RuntimeAdapter) ListMinimalAuths

func (a *RuntimeAdapter) ListMinimalAuths() []*coreauth.Auth

ListMinimalAuths returns a minimal auths.

func (*RuntimeAdapter) ListPendingPluginTasks added in v1.0.44

func (a *RuntimeAdapter) ListPendingPluginTasks(ctx context.Context, nodeType string, nodeID string) ([]node.PluginTask, error)

ListPendingPluginTasks returns plugin tasks that the node has not acked yet.

func (*RuntimeAdapter) LoadAuthIndex

func (a *RuntimeAdapter) LoadAuthIndex(ctx context.Context) error

LoadAuthIndex loads an auth index.

func (*RuntimeAdapter) LoadConfigYAML

func (a *RuntimeAdapter) LoadConfigYAML(ctx context.Context) ([]byte, error)

LoadConfigYAML loads a config yaml.

func (*RuntimeAdapter) LoadIndex

func (a *RuntimeAdapter) LoadIndex(ctx context.Context) error

LoadIndex loads an index.

func (*RuntimeAdapter) RefreshAuthIndex

func (a *RuntimeAdapter) RefreshAuthIndex(ctx context.Context, uuid string) error

RefreshAuthIndex refreshes refresh auth index.

func (*RuntimeAdapter) RemoveAuthIndex

func (a *RuntimeAdapter) RemoveAuthIndex(uuid string)

RemoveAuthIndex removes an auth index.

func (*RuntimeAdapter) ReplacePluginStatus added in v1.0.44

func (a *RuntimeAdapter) ReplacePluginStatus(ctx context.Context, nodeType string, status node.PluginTaskStatus) error

ReplacePluginStatus stores the latest plugin status report for a node.

func (*RuntimeAdapter) Save

func (a *RuntimeAdapter) Save(ctx context.Context, auth *coreauth.Auth) (string, error)

Save stores save.

func (*RuntimeAdapter) StoreAppLogPayload added in v1.0.44

func (a *RuntimeAdapter) StoreAppLogPayload(ctx context.Context, clientIP string, payload string) error

StoreAppLogPayload stores a CPA app log payload.

func (*RuntimeAdapter) StoreUsagePayload

func (a *RuntimeAdapter) StoreUsagePayload(ctx context.Context, payload string) error

StoreUsagePayload stores an usage payload.

func (*RuntimeAdapter) ValidateAPIKey added in v1.0.44

func (a *RuntimeAdapter) ValidateAPIKey(ctx context.Context, apiKey string) (bool, error)

ValidateAPIKey reports whether an API key is active in the cluster database.

type SQLiteConfig added in v1.0.44

type SQLiteConfig struct {
	Path string `yaml:"path"`
}

func (SQLiteConfig) Configured added in v1.0.44

func (c SQLiteConfig) Configured() bool

Configured reports whether SQLite has cluster database settings.

type UpsertResult added in v1.0.44

type UpsertResult string
const (
	UpsertResultCreated   UpsertResult = "created"
	UpsertResultUpdated   UpsertResult = "updated"
	UpsertResultUnchanged UpsertResult = "unchanged"
	UpsertResultRestored  UpsertResult = "restored"
)

type UsageRecord

type UsageRecord struct {
	ID                  uint      `gorm:"column:id;primaryKey;autoIncrement;index:idx_usage_time_order,priority:2"`
	Timestamp           time.Time `` /* 359-byte string literal not displayed */
	LatencyMS           int64     `gorm:"column:latency_ms;not null;default:0"`
	TTFTMS              int64     `gorm:"column:ttft_ms;not null;default:0"`
	Source              string    `gorm:"column:source;index:idx_usage_source;index:idx_usage_source_time,priority:1"`
	AuthIndex           string    `gorm:"column:auth_index;index:idx_usage_auth_index;index:idx_usage_auth_time,priority:1"`
	InputTokens         int64     `gorm:"column:input_tokens;not null;default:0"`
	OutputTokens        int64     `gorm:"column:output_tokens;not null;default:0"`
	ReasoningTokens     int64     `gorm:"column:reasoning_tokens;not null;default:0"`
	CachedTokens        int64     `gorm:"column:cached_tokens;not null;default:0"`
	CacheReadTokens     int64     `gorm:"column:cache_read_tokens;not null;default:0"`
	CacheCreationTokens int64     `gorm:"column:cache_creation_tokens;not null;default:0"`
	TotalTokens         int64     `gorm:"column:total_tokens;not null;default:0"`
	Failed              bool      `gorm:"column:failed;not null;default:false;index:idx_usage_failed;index:idx_usage_failed_time,priority:1"`
	FailStatusCode      int       `gorm:"column:fail_status_code;not null;default:0"`
	FailBody            string    `gorm:"column:fail_body;type:text"`
	Provider            string    `gorm:"column:provider;index:idx_usage_provider_model,priority:1;index:idx_usage_provider_model_time,priority:1"`
	ExecutorType        string    `gorm:"column:executor_type"`
	Model               string    `gorm:"column:model;index:idx_usage_provider_model,priority:2;index:idx_usage_provider_model_time,priority:2"`
	Alias               string    `gorm:"column:alias"`
	Effort              string    `gorm:"column:effort"`
	ServiceTier         string    `gorm:"column:service_tier"`
	Endpoint            string    `gorm:"column:endpoint;index:idx_usage_endpoint;index:idx_usage_endpoint_time,priority:1"`
	AuthType            string    `gorm:"column:auth_type"`
	APIKey              string    `gorm:"column:api_key"`
	RequestID           string    `gorm:"column:request_id;index:idx_usage_request_id"`
	HomeIP              string    `gorm:"column:home_ip;index:idx_usage_home_ip"`
	TokensJSON          JSONB     `gorm:"column:tokens"`
	FailJSON            JSONB     `gorm:"column:fail"`
	PayloadJSON         JSONB     `gorm:"column:payload;not null"`
	CreatedAt           time.Time `gorm:"column:created_at;not null"`
}

func UsageRecordFromPayload

func UsageRecordFromPayload(payload string, homeIP string) (*UsageRecord, error)

UsageRecordFromPayload derives usage record from payload.

func (UsageRecord) TableName

func (UsageRecord) TableName() string

TableName returns the database table name.

type UserRecord added in v1.0.44

type UserRecord struct {
	ID        uint           `gorm:"column:id;primaryKey;autoIncrement;index:idx_user_active_order,priority:2"`
	Username  string         `gorm:"column:username;not null;index;index:idx_user_username_active,priority:1"`
	Password  string         `gorm:"column:password;type:text"`
	Credits   float64        `gorm:"column:credits;not null;default:0"`
	MFA       JSONB          `gorm:"column:mfa"`
	Passkey   JSONB          `gorm:"column:passkey"`
	CreatedAt time.Time      `gorm:"column:created_at"`
	UpdatedAt time.Time      `gorm:"column:updated_at"`
	DeletedAt gorm.DeletedAt `gorm:"column:deleted_at;index;index:idx_user_active_order,priority:1;index:idx_user_username_active,priority:2"`
}

func (UserRecord) TableName added in v1.0.44

func (UserRecord) TableName() string

TableName returns the database table name.

type UserUpdate added in v1.0.44

type UserUpdate struct {
	Username *string
	Password *string
	Credits  *float64
	MFA      *JSONB
	Passkey  *JSONB
}

Directories

Path Synopsis

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL