Skip to content
This repository was archived by the owner on Jun 1, 2026. It is now read-only.

Commit c5e465a

Browse files
feat(api): enable redis usage queue on server
Ultraworked with [Sisyphus](https://github.com/code-yeongyu/oh-my-openagent) Co-authored-by: Sisyphus <clio-agent@sisyphuslabs.ai>
1 parent 4cbd8e3 commit c5e465a

1 file changed

Lines changed: 59 additions & 8 deletions

File tree

internal/api/server.go

Lines changed: 59 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -7,8 +7,9 @@ package api
77
import (
88
"context"
99
"crypto/subtle"
10-
"errors"
10+
"crypto/tls"
1111
"fmt"
12+
"net"
1213
"net/http"
1314
"os"
1415
"path/filepath"
@@ -27,6 +28,7 @@ import (
2728
"github.com/router-for-me/CLIProxyAPI/v6/internal/config"
2829
"github.com/router-for-me/CLIProxyAPI/v6/internal/logging"
2930
"github.com/router-for-me/CLIProxyAPI/v6/internal/managementasset"
31+
"github.com/router-for-me/CLIProxyAPI/v6/internal/redisqueue"
3032
"github.com/router-for-me/CLIProxyAPI/v6/internal/usage"
3133
"github.com/router-for-me/CLIProxyAPI/v6/internal/util"
3234
sdkaccess "github.com/router-for-me/CLIProxyAPI/v6/sdk/access"
@@ -37,6 +39,7 @@ import (
3739
sdkAuth "github.com/router-for-me/CLIProxyAPI/v6/sdk/auth"
3840
"github.com/router-for-me/CLIProxyAPI/v6/sdk/cliproxy/auth"
3941
log "github.com/sirupsen/logrus"
42+
"golang.org/x/net/http2"
4043
"gopkg.in/yaml.v3"
4144
)
4245

@@ -124,7 +127,9 @@ type Server struct {
124127
engine *gin.Engine
125128

126129
// server is the underlying HTTP server.
127-
server *http.Server
130+
server *http.Server
131+
muxBaseListener net.Listener
132+
muxHTTPListener *muxListener
128133

129134
// handlers contains the API handlers for processing requests.
130135
handlers *handlers.BaseAPIHandler
@@ -297,6 +302,7 @@ func NewServer(cfg *config.Config, authManager *auth.Manager, accessManager *sdk
297302
// or when a local management password is provided (e.g. TUI mode).
298303
hasManagementSecret := cfg.RemoteManagement.SecretKey != "" || envManagementSecret || s.localPassword != ""
299304
s.managementRoutesEnabled.Store(hasManagementSecret)
305+
redisqueue.SetEnabled(hasManagementSecret)
300306
if hasManagementSecret {
301307
s.registerManagementRoutes()
302308
}
@@ -811,23 +817,61 @@ func (s *Server) Start() error {
811817
return fmt.Errorf("failed to start HTTP server: server not initialized")
812818
}
813819

820+
baseListener, err := net.Listen("tcp", s.server.Addr)
821+
if err != nil {
822+
return fmt.Errorf("failed to listen on %s: %v", s.server.Addr, err)
823+
}
824+
s.muxBaseListener = baseListener
825+
826+
listener := net.Listener(baseListener)
814827
useTLS := s.cfg != nil && s.cfg.TLS.Enable
815828
if useTLS {
816829
cert := strings.TrimSpace(s.cfg.TLS.Cert)
817830
key := strings.TrimSpace(s.cfg.TLS.Key)
818831
if cert == "" || key == "" {
832+
_ = baseListener.Close()
819833
return fmt.Errorf("failed to start HTTPS server: tls.cert or tls.key is empty")
820834
}
835+
certPair, err := tls.LoadX509KeyPair(cert, key)
836+
if err != nil {
837+
_ = baseListener.Close()
838+
return fmt.Errorf("failed to load TLS certificate: %v", err)
839+
}
840+
tlsConfig := &tls.Config{
841+
Certificates: []tls.Certificate{certPair},
842+
NextProtos: []string{"h2", "http/1.1"},
843+
}
844+
s.server.TLSConfig = tlsConfig
845+
if err := http2.ConfigureServer(s.server, &http2.Server{}); err != nil {
846+
_ = baseListener.Close()
847+
return fmt.Errorf("failed to configure HTTP/2 server: %v", err)
848+
}
849+
listener = tls.NewListener(baseListener, tlsConfig)
821850
log.Debugf("Starting API server on %s with TLS", s.server.Addr)
822-
if errServeTLS := s.server.ListenAndServeTLS(cert, key); errServeTLS != nil && !errors.Is(errServeTLS, http.ErrServerClosed) {
823-
return fmt.Errorf("failed to start HTTPS server: %v", errServeTLS)
851+
} else {
852+
log.Debugf("Starting API server on %s", s.server.Addr)
853+
}
854+
855+
httpListener := newMuxListener(listener.Addr(), 128)
856+
s.muxHTTPListener = httpListener
857+
serveDone := make(chan error, 1)
858+
go func() {
859+
serveDone <- normalizeHTTPServeError(s.server.Serve(httpListener))
860+
}()
861+
862+
acceptErr := normalizeListenerError(s.acceptMuxConnections(listener, httpListener))
863+
if acceptErr != nil {
864+
_ = baseListener.Close()
865+
_ = httpListener.Close()
866+
_ = s.server.Close()
867+
if serveErr := <-serveDone; serveErr != nil {
868+
return fmt.Errorf("failed to start API server: accept error: %v; serve error: %v", acceptErr, serveErr)
824869
}
825-
return nil
870+
return fmt.Errorf("failed to accept API server connections: %v", acceptErr)
826871
}
827872

828-
log.Debugf("Starting API server on %s", s.server.Addr)
829-
if errServe := s.server.ListenAndServe(); errServe != nil && !errors.Is(errServe, http.ErrServerClosed) {
830-
return fmt.Errorf("failed to start HTTP server: %v", errServe)
873+
if serveErr := <-serveDone; serveErr != nil {
874+
return fmt.Errorf("failed to start API server: %v", serveErr)
831875
}
832876

833877
return nil
@@ -850,6 +894,12 @@ func (s *Server) Stop(ctx context.Context) error {
850894
default:
851895
}
852896
}
897+
if s.muxHTTPListener != nil {
898+
_ = s.muxHTTPListener.Close()
899+
}
900+
if s.muxBaseListener != nil {
901+
_ = s.muxBaseListener.Close()
902+
}
853903

854904
// Shutdown the HTTP server.
855905
if err := s.server.Shutdown(ctx); err != nil {
@@ -975,6 +1025,7 @@ func (s *Server) UpdateClients(cfg *config.Config) {
9751025
s.managementRoutesEnabled.Store(!newSecretEmpty)
9761026
}
9771027
}
1028+
redisqueue.SetEnabled(s.managementRoutesEnabled.Load())
9781029

9791030
s.applyAccessConfig(oldCfg, cfg)
9801031
s.cfg = cfg

0 commit comments

Comments
 (0)