Versions in this module Expand all Collapse all v0 v0.3.0 Jul 11, 2026 Changes in this version + type InferenceMetricsData struct + ActiveRequests int64 + AvgLatencyMs float64 + CPUUsagePercent float64 + ConnectionState string + FailedReqs int64 + GPUMetrics []GPUMetrics + LastConnectedAt time.Time + LastDisconnectedAt time.Time + MemoryTotalGB float64 + MemoryUsagePercent float64 + MemoryUsedGB float64 + ModelMetrics map[string]*ModelMetrics + P50LatencyMs float64 + P95LatencyMs float64 + P99LatencyMs float64 + ReconnectCount int64 + RequestsPerMinute float64 + StreamingReqs int64 + SuccessfulReqs int64 + TokensPerSecond float64 + TotalRequests int64 + TotalTokensIn int64 + TotalTokensOut int64 v0.2.0 Mar 26, 2026 Changes in this version + func CheckMachineIdentity(cpRepoPath string) error + type BenchmarkPayload struct + BenchmarkID string + ModelID string + Prompts []BenchmarkPrompt + TestType string + TimeoutMs int64 + type BenchmarkPrompt struct + ExpectedAnswer string + ID string + Prompt string + type BenchmarkResponsePayload struct + BenchmarkID string + Error string + RequestID string + Results []BenchmarkResult + TotalLatency int64 + type BenchmarkResult struct + Answer string + Correct bool + Error string + LatencyMs int64 + PromptID string + TokensIn int64 + TokensOut int64 type MessageType + const MsgTypeBenchmark + const MsgTypeBenchmarkResponse v0.1.1 Mar 20, 2026 v0.1.0 Mar 17, 2026 Changes in this version + const PingMsg + const PingPeriod + const TransferInterval + var ErrGlobalConcurrencyLimit = errors.New("global concurrency limit reached") + var ErrModelConcurrencyLimit = errors.New("model concurrency limit reached") + var ErrModelNotFound = &ModelError + var ErrModelQueueFull = errors.New("model queue is full") + var ErrQueueFull = errors.New("queue is full") + var ErrQueueShutdown = errors.New("queue is shutting down") + var ErrQueueStopped = errors.New("queue is stopped") + var ErrRateLimitExceeded = errors.New("rate limit exceeded") + var ErrRequestCancelled = errors.New("request was cancelled") + var ErrRequestTimeout = errors.New("request timed out") + var ErrServiceNotFound = &SupervisorError + var NetworkPolicyFlag bool + var TaskMap sync.Map + func BytesToHumanReadable(bytes int64) string + func CheckWalletBlackListForEcp(walletAddress string) bool + func CheckWalletWhiteListForEcp(walletAddress string) bool + func CronTaskForEcp() + func DoUbiTaskForDocker(c *gin.Context) + func DoZkTask(c *gin.Context) + func ExtractExposedPort(dockerfilePath string) (string, error) + func GenerateNodeID(cpRepoPath string) (string, string, string) + func GeneratePriceConfig() error + func GetAggregatedTaskInfo(taskContract string) (string, error) + func GetCpBalance() + func GetCpResource(c *gin.Context) + func GetNodeId(cpRepoPath string) string + func GetOwnerAddressAndWorkerAddress() (string, string, error) + func GetPrice(c *gin.Context) + func GetTaskInfoOnChain(taskContract string) (models.EcpTaskInfo, error) + func GetToken() (string, error) + func GetUbiResourceExporterMetrics(c *gin.Context) + func ImportImageToContainerd(tarFile string) error + func InitComputingProvider(cpRepoPath string) string + func ReceiveUbiProof(c *gin.Context) + func RestartResourceExporter() error + func RestartTraefikService() error + func RetryFn(fn func() error, maxRetries int, delay time.Duration) error + func SyncCpAccountInfo() (*models.Account, error) + func ValidateName(name string) error + type AckPayload struct + Message string + RequestID string + Success bool + type AdaptiveTokenBucket struct + func NewAdaptiveTokenBucket(initialRate float64, burstSize int, targetLatency time.Duration) *AdaptiveTokenBucket + func (atb *AdaptiveTokenBucket) RecordLatency(latency time.Duration) + type BackoffCalculator struct + func NewBackoffCalculator(strategy BackoffStrategy, initialDelay, maxDelay time.Duration) *BackoffCalculator + func (bc *BackoffCalculator) Calculate(attempt int) time.Duration + func (bc *BackoffCalculator) SetJitter(j float64) *BackoffCalculator + func (bc *BackoffCalculator) SetMultiplier(m float64) *BackoffCalculator + type BackoffStrategy int + const BackoffConstant + const BackoffExponential + const BackoffFibonacci + const BackoffLinear + type CircuitBreaker struct + func NewCircuitBreaker(failureThreshold, successThreshold int, timeout time.Duration) *CircuitBreaker + func (cb *CircuitBreaker) Allow() bool + func (cb *CircuitBreaker) GetState() CircuitState + func (cb *CircuitBreaker) RecordFailure() + func (cb *CircuitBreaker) RecordSuccess() + func (cb *CircuitBreaker) Reset() + func (cb *CircuitBreaker) SetStateChangeCallback(callback func(from, to CircuitState)) + type CircuitState int + const CircuitClosed + const CircuitHalfOpen + const CircuitOpen + func (s CircuitState) String() string + type ConcurrencyConfig struct + AcquireTimeout time.Duration + DefaultModelMax int + EnableGPUAwareness bool + GPUMemoryBufferMB int + GlobalMaxConcurrent int + func DefaultConcurrencyConfig() ConcurrencyConfig + type ConcurrencyLimiter struct + func NewConcurrencyLimiter(config ConcurrencyConfig, gpuCollector *GPUMetricsCollector) *ConcurrencyLimiter + func (cl *ConcurrencyLimiter) Acquire(ctx context.Context, modelID string) (*ConcurrencyToken, error) + func (cl *ConcurrencyLimiter) GetMetrics() ConcurrencyMetrics + func (cl *ConcurrencyLimiter) GetModelConcurrency(modelID string) (current, max int) + func (cl *ConcurrencyLimiter) RegisterModel(modelID string, maxConcurrent int, gpuMemoryMB int) + func (cl *ConcurrencyLimiter) SetGlobalMax(max int) + func (cl *ConcurrencyLimiter) SetModelMax(modelID string, max int) + func (cl *ConcurrencyLimiter) Start() + func (cl *ConcurrencyLimiter) Stop() + func (cl *ConcurrencyLimiter) TryAcquire(modelID string) (*ConcurrencyToken, error) + func (cl *ConcurrencyLimiter) UnregisterModel(modelID string) + type ConcurrencyMetrics struct + AvgHoldTimeMs float64 + GlobalActive int64 + GlobalMax int + PerModelActive map[string]int64 + PerModelMax map[string]int + TotalAcquired int64 + TotalRejected int64 + TotalReleased int64 + TotalTimeouts int64 + type ConcurrencyToken struct + func (ct *ConcurrencyToken) Release() + type Config struct + Resources map[string]string + type CpBalanceService struct + func NewCpBalanceService() CpBalanceService + func (cpServ CpBalanceService) GetCpBalance(cpAccount string) (*models.CpBalanceEntity, error) + func (cpServ CpBalanceService) SaveCpBalance(cpBalance models.CpBalanceEntity) (err error) + func (cpServ CpBalanceService) UpdateCpBalance(cpBalance models.CpBalanceEntity) error + type CpInfoService struct + func NewCpInfoService() CpInfoService + func (cpServ CpInfoService) GetCpInfoEntityByAccountAddress(accountAddress string) (*models.CpInfoEntity, error) + func (cpServ CpInfoService) SaveCpInfoEntity(cp *models.CpInfoEntity) (err error) + func (cpServ CpInfoService) UpdateCpInfoByNodeId(cp *models.CpInfoEntity) (err error) + type CronTask struct + func NewCronTask(nodeId string) *CronTask + func (task *CronTask) CheckCpBalance() + func (task *CronTask) DeleteSpaceLog() + func (task *CronTask) RunTask() + type DeterministicChallengeData struct + MaxTokens int + Prompt string + Seed int + type DeterministicResponseData struct + Text string + Tokens []string + type DockerService struct + func NewDockerService() *DockerService + func (ds *DockerService) BuildImage(jobUuid, buildPath, imageName string) error + func (ds *DockerService) CheckRunningContainer(containerName string) (bool, string, error) + func (ds *DockerService) CleanResourceForDocker(onlyClearContainer bool) + func (ds *DockerService) ContainerCreateAndStart(config *container.Config, hostConfig *container.HostConfig, ...) error + func (ds *DockerService) ContainerExec(containerID string, cmd []string) (string, error) + func (ds *DockerService) ContainerLogs(containerName string) (string, error) + func (ds *DockerService) CreateNetwork(networkName string) error + func (ds *DockerService) GetContainerLogStream(ctx context.Context, containerName string) (io.ReadCloser, error) + func (ds *DockerService) GetContainerStatus() (map[string]string, error) + func (ds *DockerService) IsExistContainer(containerName string) bool + func (ds *DockerService) PullImage(imageName string) error + func (ds *DockerService) PushImage(imagesName string) error + func (ds *DockerService) RemoveContainerByName(containerName string) error + func (ds *DockerService) RemoveImage(imageId string) error + func (ds *DockerService) SaveDockerImage(imageName string) (string, error) + type EcpJobService struct + func NewEcpJobService() EcpJobService + func (cpServ EcpJobService) DeleteContainerByUuid(uuid string) (err error) + func (cpServ EcpJobService) GetEcpJobByUuid(uuid string) (*models.EcpJobEntity, error) + func (cpServ EcpJobService) GetEcpJobList(status []string) ([]models.EcpJobEntity, error) + func (cpServ EcpJobService) GetEcpJobs(jobUuid string) ([]models.EcpJobEntity, error) + func (cpServ EcpJobService) GetEcpJobsByLimit(tailNum int) ([]models.EcpJobEntity, error) + func (cpServ EcpJobService) SaveEcpJobEntity(job *models.EcpJobEntity) (err error) + func (cpServ EcpJobService) UpdateEcpJobEntity(jobUuid, status string) (err error) + func (cpServ EcpJobService) UpdateEcpJobEntityContainerName(jobUuid string, containerName string) (err error) + func (cpServ EcpJobService) UpdateEcpJobEntityMessage(jobUuid string, message string) (err error) + func (cpServ EcpJobService) UpdateEcpJobEntityPortsAndServiceUrl(jobUuid, portMap, serviceUrl string) (err error) + func (cpServ EcpJobService) UpdateEcpJobEntityRewardAndBlock(jobUuid string, blockNumber int64, reward float64) (err error) + type ErrorLine struct + Error string + ErrorDetail struct{ ... } + type ErrorPayload struct + Code int + Message string + RequestID string + type FingerprintChallengeData struct + Files []FingerprintChallengeFile + type FingerprintChallengeFile struct + ExpectedHash string + Filename string + type FingerprintResponseData struct + Files []FingerprintResponseFile + type FingerprintResponseFile struct + Filename string + Hash string + Status string + type GPUInfo struct + Available int + FreeIndex []string + Name string + Total int + Used int + type GPUManager struct + func NewGpuManager() *GPUManager + func (gm *GPUManager) AllocateGPU(key string, count int) ([]string, error) + func (gm *GPUManager) CheckAvailableGPU() (string, bool) + func (gm *GPUManager) GetGPU(key string) (GPUInfo, bool) + func (gm *GPUManager) ReleaseGPU(key string, count int, indexs []string) error + func (gm *GPUManager) UpdateGPU(key string, gpuInfo GPUInfo) + type GPUMetrics struct + ComputeProcesses int + FanSpeedPct float64 + Index int + MemoryTotalMB float64 + MemoryUsagePct float64 + MemoryUsedMB float64 + Name string + PowerDrawW float64 + PowerLimitW float64 + TemperatureC float64 + UUID string + UtilizationPct float64 + type GPUMetricsCollector struct + func NewGPUMetricsCollector() *GPUMetricsCollector + func (c *GPUMetricsCollector) CollectGPUMetrics() []GPUMetrics + func (c *GPUMetricsCollector) GetAggregatedGPUMetrics() (avgUtilization, avgMemoryUsage float64) + func (c *GPUMetricsCollector) IsGPUAvailable() bool + type HardwareField struct + Name string + TagValue int + Value string + func GetStructByTag(v interface{}) ([]HardwareField, error) + type HardwareInfo struct + CUDAVersion string + ComputeCapability string + DriverVersion string + GPUCount int + GPUModel string + GPUType string + ServingEngine string + VRAMGB int + func DetectGPUHardware() *HardwareInfo + type HardwarePrice struct + GpusPrice map[string]string + TARGET_CPU string + TARGET_GPU_DEFAULT string + TARGET_HD_EPHEMERAL string + TARGET_MEMORY string + func ReadPriceConfig() (HardwarePrice, error) + type HealthCheckConfig struct + CircuitOpenTime time.Duration + HealthyThreshold int + Interval time.Duration + Timeout time.Duration + UnhealthyThreshold int + func DefaultHealthCheckConfig() HealthCheckConfig + type HeartbeatPayload struct + Hardware *HardwareInfo + Metrics map[string]float64 + ModelHealth map[string]string + Models []string + NodeID string + ProviderID string + Timestamp int64 + type HistoricalDataPoint struct + AvgLatencyMs float64 + P99LatencyMs float64 + RequestsPerMinute float64 + SuccessRate float64 + Timestamp time.Time + TokensPerSecond float64 + TotalRequests int64 + type HttpClient struct + func NewHttpClient(host string, header http.Header) *HttpClient + func (c *HttpClient) Get(api string, queries url.Values, dest any) error + func (c *HttpClient) PostForm(api string, data url.Values, dest any) error + func (c *HttpClient) PostJSON(api string, data any, dest any) error + func (c *HttpClient) Request(method string, api string, body io.Reader, dest any, contentType ...string) (err error) + type ImageJobService struct + func NewImageJobService() *ImageJobService + func (*ImageJobService) CheckJobCondition(c *gin.Context) + func (*ImageJobService) DeleteJob(c *gin.Context) + func (*ImageJobService) DeployInference(c *gin.Context, deployJob models.DeployJobParam, totalCost float64, ...) + func (*ImageJobService) DeployMining(c *gin.Context, deployJob models.DeployJobParam, totalCost float64, ...) + func (*ImageJobService) DockerLogsHandler(c *gin.Context) + func (*ImageJobService) GetJobStatus(c *gin.Context) + func (imageJob *ImageJobService) DeployJob(c *gin.Context) + type InferenceClient struct + func NewInferenceClient(nodeID, workerAddr, ownerAddr string) *InferenceClient + func (c *InferenceClient) GetMetrics() InferenceMetrics + func (c *InferenceClient) GetMetricsPrometheus() string + func (c *InferenceClient) GetNodeID() string + func (c *InferenceClient) IsConnected() bool + func (c *InferenceClient) SendModelHealthUpdate(modelHealth map[string]string) + func (c *InferenceClient) SetInferenceHandler(handler InferenceHandler) + func (c *InferenceClient) SetModelHealthProvider(provider func() map[string]string) + func (c *InferenceClient) SetModelMappingsProvider(provider func() map[string]ModelMapping) + func (c *InferenceClient) SetStreamingInferenceHandler(handler StreamingInferenceHandler) + func (c *InferenceClient) SetWarmupHandler(handler WarmupHandler) + func (c *InferenceClient) Start() error + func (c *InferenceClient) Stop() + type InferenceHandler func(payload InferencePayload) (*InferenceResponse, error) + type InferenceMetrics struct + ActiveRequests int64 + AvgLatencyMs float64 + CPUUsagePercent float64 + ConnectionState string + FailedReqs int64 + GPUMetrics []GPUMetrics + LastConnectedAt time.Time + LastDisconnectedAt time.Time + MemoryTotalGB float64 + MemoryUsagePercent float64 + MemoryUsedGB float64 + ModelMetrics map[string]*ModelMetrics + P50LatencyMs float64 + P95LatencyMs float64 + P99LatencyMs float64 + ReconnectCount int64 + RequestsPerMinute float64 + StreamingReqs int64 + SuccessfulReqs int64 + TokensPerSecond float64 + TotalRequests int64 + TotalTokensIn int64 + TotalTokensOut int64 + func NewInferenceMetrics() *InferenceMetrics + func (m *InferenceMetrics) GetPrometheusMetrics() string + func (m *InferenceMetrics) GetRequestHistory(limit int, modelFilter string) []RequestMetric + func (m *InferenceMetrics) GetSnapshot() InferenceMetrics + func (m *InferenceMetrics) RecordConnectionState(state string) + func (m *InferenceMetrics) RecordReconnect() + func (m *InferenceMetrics) RecordRequest(req RequestMetric) + func (m *InferenceMetrics) RecordRequestEnd(model string, latencyMs float64, tokensIn, tokensOut int, success bool, ...) + func (m *InferenceMetrics) RecordRequestStart(model string, streaming bool) + func (m *InferenceMetrics) Reset() + func (m *InferenceMetrics) UpdateGPUMetrics(gpuMetrics []GPUMetrics) + func (m *InferenceMetrics) UpdateSystemMetrics(cpuPercent, memPercent, memUsedGB, memTotalGB float64) + type InferencePayload struct + EndpointID string + ModelID string + Request json.RawMessage + Stream bool + type InferenceResponse struct + Error string + Latency int64 + RequestID string + Response json.RawMessage + StatusCode int + type InferenceService struct + func NewInferenceService(nodeID, cpPath string) *InferenceService + func (s *InferenceService) DisableModel(modelID string) error + func (s *InferenceService) EnableModel(modelID string) error + func (s *InferenceService) ForceHealthCheck(modelID string) + func (s *InferenceService) GetActiveModels() []string + func (s *InferenceService) GetAllModelHealth() map[string]*ModelStatus + func (s *InferenceService) GetAllModels() []*RegisteredModel + func (s *InferenceService) GetClient() *InferenceClient + func (s *InferenceService) GetConcurrencyMetrics() *ConcurrencyMetrics + func (s *InferenceService) GetHealthChecker() *ModelHealthChecker + func (s *InferenceService) GetMetrics() *InferenceMetrics + func (s *InferenceService) GetMetricsHistory(duration, resolution time.Duration) ([]HistoricalDataPoint, error) + func (s *InferenceService) GetMetricsPrometheus() string + func (s *InferenceService) GetModelDetailedMetrics(modelID string) map[string]interface{} + func (s *InferenceService) GetModelHealth(modelID string) (*ModelStatus, bool) + func (s *InferenceService) GetModelStatus(modelID string) (*RegisteredModel, bool) + func (s *InferenceService) GetModelsSummary() map[string]interface{} + func (s *InferenceService) GetRateLimiterMetrics() *RateLimiterMetrics + func (s *InferenceService) GetRegistry() *ModelRegistry + func (s *InferenceService) GetRequestHistory(limit int, modelFilter string) []RequestMetric + func (s *InferenceService) GetRequestManagementStatus() map[string]interface{} + func (s *InferenceService) GetRetryMetrics() *RetryMetrics + func (s *InferenceService) IsConnected() bool + func (s *InferenceService) IsHealthy() bool + func (s *InferenceService) Name() string + func (s *InferenceService) RegisterModels(models []string) + func (s *InferenceService) ReloadModels() error + func (s *InferenceService) SetGlobalConcurrencyLimit(max int) + func (s *InferenceService) SetGlobalRateLimit(tokensPerSecond float64) + func (s *InferenceService) SetModelConcurrencyLimit(modelID string, max int) + func (s *InferenceService) SetModelRateLimit(modelID string, tokensPerSecond float64, burstSize int) + func (s *InferenceService) Start() error + func (s *InferenceService) Stop() + type JobService struct + func NewJobService() JobService + func (jobServ JobService) DeleteJobEntityByJobUuId(jobUuid string, jobStatus int) error + func (jobServ JobService) DeleteJobEntityBySpaceUuId(spaceUuid, jobUuid string, jobStatus int) error + func (jobServ JobService) GetJobEntityByJobUuid(jobUuid string) (models.JobEntity, error) + func (jobServ JobService) GetJobEntityBySpaceUuid(spaceUuid string) int64 + func (jobServ JobService) GetJobEntityByTaskUuid(taskUuid string) (models.JobEntity, error) + func (jobServ JobService) GetJobList(status int, tailNum int) (list []*models.JobEntity, err error) + func (jobServ JobService) GetJobListByNoRejectStatus() (list []*models.JobEntity, err error) + func (jobServ JobService) GetJobListByNoReward() (list []*models.JobEntity, err error) + func (jobServ JobService) SaveJobEntity(job *models.JobEntity) (err error) + func (jobServ JobService) UpdateJobEntityByJobUuid(job *models.JobEntity) (err error) + func (jobServ JobService) UpdateJobEntityStatusByJobUuid(jobUuid string, status int) (err error) + func (jobServ JobService) UpdateJobResultUrlByJobUuid(jobUuid string, resultUrl string) (err error) + func (jobServ JobService) UpdateJobReward(taskUuid string, amount string) (err error) + func (jobServ JobService) UpdateJobScannedBlock(taskUuid string, end uint64) (err error) + type Message struct + Payload json.RawMessage + RequestID string + Type MessageType + type MessageType string + const MsgTypeAck + const MsgTypeError + const MsgTypeHeartbeat + const MsgTypeInference + const MsgTypeModelHealthUpdate + const MsgTypeRegister + const MsgTypeStreamChunk + const MsgTypeStreamEnd + const MsgTypeVerify + const MsgTypeWarmup + type MetricsHistory struct + func NewMetricsHistory() *MetricsHistory + func (h *MetricsHistory) GetHistory(duration time.Duration, resolution time.Duration) ([]HistoricalDataPoint, error) + func (h *MetricsHistory) GetRecentDataPoints(count int) ([]HistoricalDataPoint, error) + func (h *MetricsHistory) Start(metricsProvider func() *InferenceMetrics) error + func (h *MetricsHistory) Stop() + type MetricsHistoryEntity struct + ActiveRequests int64 + AvgLatencyMs float64 + FailedReqs int64 + ID uint + P50LatencyMs float64 + P95LatencyMs float64 + P99LatencyMs float64 + RequestsPerMinute float64 + SuccessRate float64 + SuccessfulReqs int64 + Timestamp time.Time + TokensPerSecond float64 + TotalRequests int64 + TotalTokensIn int64 + TotalTokensOut int64 + func (MetricsHistoryEntity) TableName() string + type ModelError struct + Message string + func (e *ModelError) Error() string + type ModelHealth int + const ModelHealthDegraded + const ModelHealthHealthy + const ModelHealthUnhealthy + const ModelHealthUnknown + func (h ModelHealth) String() string + type ModelHealthChecker struct + func NewModelHealthChecker(config HealthCheckConfig) *ModelHealthChecker + func (h *ModelHealthChecker) ForceCheck(modelID string) + func (h *ModelHealthChecker) GetAllStatuses() map[string]*ModelStatus + func (h *ModelHealthChecker) GetHealthyModels() []string + func (h *ModelHealthChecker) GetModelStatus(modelID string) (*ModelStatus, bool) + func (h *ModelHealthChecker) GetStatusJSON() ([]byte, error) + func (h *ModelHealthChecker) IsModelHealthy(modelID string) bool + func (h *ModelHealthChecker) RegisterModel(modelID, endpoint, apiKey string) + func (h *ModelHealthChecker) SetStatusChangeCallback(cb func(modelID string, oldHealth, newHealth ModelHealth)) + func (h *ModelHealthChecker) Start() + func (h *ModelHealthChecker) Stop() + func (h *ModelHealthChecker) UnregisterModel(modelID string) + type ModelHealthUpdatePayload struct + ModelHealth map[string]string + NodeID string + ProviderID string + Timestamp int64 + type ModelInfo struct + Format string + HashAlgo string + ModelID string + Quantization string + WeightHash string + type ModelMapping struct + APIKey string + Category string + Container string + Endpoint string + Format string + GPUMemory int + LocalModel string + Quantization string + type ModelMetrics struct + ActiveRequests int64 + AvgLatencyMs float64 + FailedReqs int64 + ModelName string + SuccessfulReqs int64 + TokensPerSecond float64 + TotalRequests int64 + TotalTokensIn int64 + TotalTokensOut int64 + type ModelRegistry struct + func NewModelRegistry(configPath string, healthChecker *ModelHealthChecker) *ModelRegistry + func (r *ModelRegistry) DisableModel(modelID string) error + func (r *ModelRegistry) EnableModel(modelID string) error + func (r *ModelRegistry) GetAllModelHealthMap() map[string]string + func (r *ModelRegistry) GetAllModels() []*RegisteredModel + func (r *ModelRegistry) GetLocalModelName(modelID string) string + func (r *ModelRegistry) GetModel(modelID string) (*RegisteredModel, bool) + func (r *ModelRegistry) GetModelAPIKey(modelID string) string + func (r *ModelRegistry) GetModelEndpoint(modelID string) (string, bool) + func (r *ModelRegistry) GetModelMappings() map[string]ModelMapping + func (r *ModelRegistry) GetReadyModelIDs() []string + func (r *ModelRegistry) GetReadyModels() []*RegisteredModel + func (r *ModelRegistry) GetStatusSummary() map[string]interface{} + func (r *ModelRegistry) ReloadConfig() error + func (r *ModelRegistry) SetCallbacks(onAdded func(model *RegisteredModel), onRemoved func(modelID string), ...) + func (r *ModelRegistry) SetHealthUpdateCallback(callback func(modelHealth map[string]string)) + func (r *ModelRegistry) Start() error + func (r *ModelRegistry) Stop() + type ModelServerError struct + Body []byte + Message string + StatusCode int + func (e *ModelServerError) Error() string + type ModelState int + const ModelStateDisabled + const ModelStateLoading + const ModelStateReady + const ModelStateUnhealthy + const ModelStateUnknown + func (s ModelState) String() string + type ModelStatus struct + AvgLatencyMs float64 + CircuitOpen bool + ConsecutiveFails int + Endpoint string + Health ModelHealth + HealthString string + LastCheck time.Time + LastError string + LastSuccess time.Time + LatencyMs float64 + ModelID string + TotalChecks int64 + TotalFailures int64 + TotalSuccesses int64 + type PerModelRateLimiter struct + func NewPerModelRateLimiter(config RateLimiterConfig, gpuCollector *GPUMetricsCollector, ...) *PerModelRateLimiter + func (prl *PerModelRateLimiter) EnsureModelLimit(modelID string) + type ProviderStatusResponse struct + APIKeyValid bool + CanConnect bool + EarningsEnabled bool + Message string + Name string + NextSteps []string + ProviderID string + Status string + Step int + StepLabel string + TotalSteps int + Warning string + type QueueConfig struct + DefaultTimeout time.Duration + DrainTimeout time.Duration + EnablePriority bool + MaxPerModelQueue int + MaxQueueSize int + func DefaultQueueConfig() QueueConfig + type QueueMetrics struct + AvgWaitTimeMs float64 + CurrentDepth int64 + MaxWaitTimeMs float64 + PerModelDepth map[string]int64 + TotalCancelled int64 + TotalDequeued int64 + TotalEnqueued int64 + TotalRejected int64 + TotalTimedOut int64 + type QueueResult struct + Error error + Response json.RawMessage + type QueuedRequest struct + Deadline time.Time + EnqueueTime time.Time + ID string + ModelID string + Payload json.RawMessage + Priority RequestPriority + ResultChan chan *QueueResult + type RateLimiter struct + func NewRateLimiter(config RateLimiterConfig, gpuCollector *GPUMetricsCollector) *RateLimiter + func (rl *RateLimiter) Allow() bool + func (rl *RateLimiter) AllowModel(modelID string) bool + func (rl *RateLimiter) GetBackoffTime() time.Duration + func (rl *RateLimiter) GetMetrics() RateLimiterMetrics + func (rl *RateLimiter) GetModelMetrics(modelID string) *RateLimiterMetrics + func (rl *RateLimiter) RemoveModelLimit(modelID string) + func (rl *RateLimiter) SetModelLimit(modelID string, tokensPerSecond float64, burstSize int) + func (rl *RateLimiter) Start() + func (rl *RateLimiter) Stop() + func (rl *RateLimiter) WaitForToken(modelID string, timeout time.Duration) error + type RateLimiterConfig struct + AdaptiveAdjustment float64 + AdaptiveMaxRate float64 + AdaptiveMinRate float64 + BurstSize int + EnableAdaptive bool + GPUThresholdHigh float64 + GPUThresholdLow float64 + TokensPerSecond float64 + func DefaultRateLimiterConfig() RateLimiterConfig + type RateLimiterMetrics struct + AdaptiveEnabled bool + BurstSize int + CurrentRate float64 + CurrentTokens float64 + TotalAllowed int64 + TotalThrottled int64 + type RegisterPayload struct + Capabilities []string + Hardware *HardwareInfo + ModelHashes []ModelInfo + Models []string + NodeID string + NodeName string + OwnerAddr string + ProviderID string + Signature string + Token string + WorkerAddr string + type RegisteredModel struct + APIKey string + Category string + Container string + Enabled bool + Endpoint string + Format string + GPUMemory int + Health ModelHealth + HealthString string + ID string + LoadedAt time.Time + LocalModel string + Quantization string + State ModelState + StateString string + UpdatedAt time.Time + type RequestMetric struct + EndTime time.Time + ErrorReason string + LatencyMs float64 + Model string + RequestID string + StartTime time.Time + Streaming bool + Success bool + TokensIn int + TokensOut int + type RequestPriority int + const PriorityHigh + const PriorityLow + const PriorityNormal + type RequestQueue struct + func NewRequestQueue(config QueueConfig) *RequestQueue + func (q *RequestQueue) Dequeue() *QueuedRequest + func (q *RequestQueue) Enqueue(req *QueuedRequest) error + func (q *RequestQueue) EnqueueWithTimeout(req *QueuedRequest, timeout time.Duration) (*QueueResult, error) + func (q *RequestQueue) GetDepth() int + func (q *RequestQueue) GetMetrics() QueueMetrics + func (q *RequestQueue) GetModelDepth(modelID string) int + func (q *RequestQueue) IsAccepting() bool + func (q *RequestQueue) SetProcessFunc(f func(*QueuedRequest)) + func (q *RequestQueue) Start() + func (q *RequestQueue) Stop() + type RestartableService interface + IsHealthy func() bool + Name func() string + Start func() error + Stop func() + type ResultChecker interface + Check func() error + type RetryConfig struct + InitialDelay time.Duration + JitterFactor float64 + MaxDelay time.Duration + MaxRetries int + Multiplier float64 + NonRetryableErrors []string + RetryableErrors []string + func DefaultRetryConfig() RetryConfig + type RetryMetrics struct + AvgRetriesPerRequest float64 + RetrySuccessRate float64 + TotalAttempts int64 + TotalFailures int64 + TotalNonRetryable int64 + TotalRetries int64 + TotalSuccesses int64 + type RetryPolicy struct + func NewRetryPolicy(config RetryConfig) *RetryPolicy + func (rp *RetryPolicy) CalculateDelay(attempt int) time.Duration + func (rp *RetryPolicy) Execute(ctx context.Context, operation func() error) error + func (rp *RetryPolicy) ExecuteWithResult(ctx context.Context, operation func() (interface{}, error)) error + func (rp *RetryPolicy) GetMetrics() RetryMetrics + func (rp *RetryPolicy) IsRetryable(err error) bool + type RetryableOperation struct + func NewRetryableOperation(policy *RetryPolicy, op func(ctx context.Context) error) *RetryableOperation + func (ro *RetryableOperation) OnRetry(callback func(attempt int, err error)) *RetryableOperation + func (ro *RetryableOperation) Run(ctx context.Context) error + type Semaphore struct + func NewSemaphore(max int) *Semaphore + func (s *Semaphore) Acquire(timeout time.Duration) bool + func (s *Semaphore) GetStats() (current, max int, acquired, released, rejected, timeouts int64, ...) + func (s *Semaphore) Release(holdTime time.Duration) + func (s *Semaphore) SetMax(max int) + func (s *Semaphore) TryAcquire() bool + type SendProofResp struct + Code int + Data struct{ ... } + Msg string + type SequenceTask struct + CheckCode string + Deadline int + Id int + InputParam string + Proof string + ResourceType int + Reward string + SequenceCid string + SequenceTaskAddr string + SettlementCid string + SettlementTaskAddr string + Status string + Type int + Uuid string + VerifyParam string + type Sequencer struct + func NewSequencer() *Sequencer + func (s *Sequencer) GetToken() error + func (s *Sequencer) QueryTask(taskType int, taskIds []int64, uuids []string) (TaskListResp, error) + func (s *Sequencer) SendTaskProof(data []byte) (SendProofResp, error) + type ServiceState struct + CurrentBackoff time.Duration + Healthy bool + LastHealthyAt time.Time + LastRestartAt time.Time + Name string + RestartCount int + type StreamChunkPayload struct + Chunk json.RawMessage + Done bool + RequestID string + type StreamEndPayload struct + Error string + Latency int64 + RequestID string + StatusCode int + TokensInput int64 + TokensOutput int64 + type StreamResult struct + Error error + TokensInput int64 + TokensOutput int64 + type StreamingInferenceHandler func(requestID string, payload InferencePayload, ...) *StreamResult + type Supervisor struct + func NewSupervisor(config SupervisorConfig) *Supervisor + func (s *Supervisor) ForceRestart(name string) error + func (s *Supervisor) GetServiceState(name string) (ServiceState, bool) + func (s *Supervisor) GetServiceStates() map[string]ServiceState + func (s *Supervisor) Register(service RestartableService) + func (s *Supervisor) Start() + func (s *Supervisor) Stop() + func (s *Supervisor) Unregister(name string) + type SupervisorConfig struct + HealthCheckInterval time.Duration + MaxRestartAttempts int + MaxRestartBackoff time.Duration + RestartBackoff time.Duration + func DefaultSupervisorConfig() SupervisorConfig + type SupervisorError struct + Message string + func (e *SupervisorError) Error() string + type TaskGroup struct + Ids []int64 + Items []*models.TaskEntity + Type int + Uuids []string + type TaskListResp struct + Code int + Data struct{ ... } + Msg string + type TaskPaymentService struct + func NewTaskPaymentService() *TaskPaymentService + func (tps *TaskPaymentService) ScannerChainGetTaskPayment() + type TaskService struct + func NewTaskService() TaskService + func (taskServ TaskService) GetTaskByUuid(uuid string) (*models.TaskEntity, error) + func (taskServ TaskService) GetTaskEntity(taskId int64) (*models.TaskEntity, error) + func (taskServ TaskService) GetTaskList(tailNum int, taskStatus ...int) (list []*models.TaskEntity, err error) + func (taskServ TaskService) GetTaskListNoRewardForFilC2() (list []*models.TaskEntity, err error) + func (taskServ TaskService) GetTaskListNoRewardForMining() (list []*models.TaskEntity, err error) + func (taskServ TaskService) SaveTaskEntity(task *models.TaskEntity) (err error) + func (taskServ TaskService) UpdateTaskEntityByTaskId(task *models.TaskEntity) (err error) + func (taskServ TaskService) UpdateTaskEntityByTaskUuId(task *models.TaskEntity) (err error) + func (taskServ TaskService) UpdateTaskStatusById(taskId int, status int) (err error) + func (taskServ TaskService) UpdateTaskStatusByUuid(uuid string, status int) (err error) + type TaskStatus struct + Code int + Data struct{ ... } + Msg string + type TokenBucket struct + func NewTokenBucket(tokensPerSecond float64, burstSize int) *TokenBucket + func (tb *TokenBucket) Allow() bool + func (tb *TokenBucket) AllowN(n int) bool + func (tb *TokenBucket) GetStats() (tokens float64, rate float64, allowed, throttled int64) + func (tb *TokenBucket) SetRate(tokensPerSecond float64) + type TokenResp struct + Code int + Data struct{ ... } + Msg string + type VerifyPayload struct + Challenge json.RawMessage + ChallengeID string + ChallengeType string + ModelID string + type VerifyResponsePayload struct + ChallengeID string + Error string + Response json.RawMessage + Success bool + type WarmupHandler func(payload WarmupPayload) (*WarmupResponse, error) + type WarmupPayload struct + ModelID string + WarmupType string + type WarmupResponse struct + Error string + LoadTimeMs int64 + MemoryMB int64 + ModelID string + RequestID string + Success bool + type WsClient struct + func NewWsClient(client *websocket.Conn) *WsClient + func (ws *WsClient) Close() + func (ws *WsClient) HandleLogs(reader io.Reader) + type YamlStruct struct + Services struct{ ... }