Files
pangolin/server/cmd/server/main.go
T
wangjia 4940278ea5
Deploy Server / deploy-server (push) Successful in 2m55s
Deploy Site / deploy-site (push) Successful in 2m41s
Deploy Client / build-windows (push) Successful in 1m47s
Deploy Client / build-android (push) Successful in 7m39s
Deploy Client / build-macos (push) Successful in 3m40s
Deploy Client / build-ios (push) Successful in 4m48s
Deploy Client / release-deploy (push) Successful in 2m25s
feat(migrate): 用户中心迁到 pangolin.yanmeiai.com/user/ + 域名配置化
原独立子域 app.yanmeiai.com → 主站子路径 /user/(用户选停用旧域名):
- 域名配置化(去硬编码):客户端 kWebUserCenterBaseUrl 收进 api_config.dart
  (dart-define 可覆盖);官网 site.ts、服务端 CORS_ORIGINS 本就是配置
- 用户中心 next.config basePath=/user;layout.tsx 的 /colors_and_type.css 手动
  拼 basePath(public 根绝对资源不自动加前缀,否则 404)
- CI 合并部署:compile-site + compile-usercenter + combine-site(用户中心并入
  dist/user/ + _headers 按 /user/* 分域:官网严格 CSP,用户中心 unsafe-inline+
  connect-src https)→ 单次部署 pangolin-site;删独立 pangolin-usercenter 部署
- 服务端 CORS 默认 origin app.yanmeiai.com → pangolin.yanmeiai.com(+ 测试)
- 客户端 web_launch 走新址 → 随 client-v1.0.62;go test CORS 过、flutter analyze 净

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_013nMthbVEmQquxBRKb9Fj8u
2026-07-07 17:33:22 +08:00

598 lines
24 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
package main
import (
"context"
"crypto/tls"
"database/sql"
"encoding/hex"
"encoding/json"
"flag"
"log"
"log/slog"
"net"
"net/http"
"os"
"os/signal"
"strconv"
"strings"
"syscall"
"time"
"github.com/go-chi/chi/v5"
chimw "github.com/go-chi/chi/v5/middleware"
"github.com/redis/go-redis/v9"
"google.golang.org/grpc"
"google.golang.org/grpc/credentials"
"github.com/wangjia/pangolin/server/internal/admin"
"github.com/wangjia/pangolin/server/internal/apierr"
"github.com/wangjia/pangolin/server/internal/auth"
"github.com/wangjia/pangolin/server/internal/codes"
"github.com/wangjia/pangolin/server/internal/db"
"github.com/wangjia/pangolin/server/internal/devices"
"github.com/wangjia/pangolin/server/internal/httpapi"
"github.com/wangjia/pangolin/server/internal/mtls"
"github.com/wangjia/pangolin/server/internal/nodes"
agentv1 "github.com/wangjia/pangolin/server/internal/pb/agentv1"
"github.com/wangjia/pangolin/server/internal/provision"
"github.com/wangjia/pangolin/server/internal/provision/providers"
"github.com/wangjia/pangolin/server/internal/redisutil"
"github.com/wangjia/pangolin/server/internal/scheduler"
"github.com/wangjia/pangolin/server/internal/scheduler/probe"
"github.com/wangjia/pangolin/server/internal/sessions"
"github.com/wangjia/pangolin/server/internal/usage"
)
func main() {
listenAddr := flag.String("addr", "", "HTTP listen address (default :8080, overridden by ADDR env)")
flag.Parse()
if *listenAddr == "" {
if v := os.Getenv("ADDR"); v != "" {
*listenAddr = v
} else {
*listenAddr = ":8080"
}
}
// Admin backend (separate internal listener, optional).
startAdminIfConfigured()
// ─── Shared: Redis ────────────────────────────────────────────────────────
redisAddr := getenvDefault("REDIS_ADDR", "127.0.0.1:6379")
rdb, err := redisutil.New(redisAddr, os.Getenv("REDIS_PASSWORD"), 0)
if err != nil {
log.Fatalf("redis connect: %v", err)
}
// ─── Shared: DB ───────────────────────────────────────────────────────────
dbDSN := os.Getenv("DB_DSN")
var sqlDB *sql.DB
if dbDSN != "" {
sqlDB, err = db.Open(dbDSN)
if err != nil {
log.Fatalf("db open: %v", err)
}
}
// ─── Shared: nodes.Service (Hub shared by HTTP connect + gRPC agent) ──────
// Constructed here so both sides use the same Hub instance.
var nodeSvc *nodes.Service
if sqlDB != nil {
nodeStore := nodes.NewSQLNodeStore(sqlDB)
grpcAddr := os.Getenv("GRPC_ADDR")
caKeyPath := os.Getenv("CA_KEY_PATH")
caCertPath := os.Getenv("CA_CERT_PATH")
grpcCertPath := os.Getenv("GRPC_CERT_PATH")
grpcKeyPath := os.Getenv("GRPC_KEY_PATH")
if grpcAddr != "" && caKeyPath != "" && caCertPath != "" &&
grpcCertPath != "" && grpcKeyPath != "" {
ca, err := mtls.NewCA(mtls.CAConfig{KeyPath: caKeyPath, CertPath: caCertPath})
if err != nil {
log.Fatalf("load node CA: %v", err)
}
tokens := mtls.NewBootstrapTokenManager(rdb)
crl := mtls.NewCRL(rdb, nil)
nodeSvc = nodes.NewService(ca, tokens, rdb, nodeStore)
nodeSvc.Hub().Start(context.Background())
go startGRPCWithService(nodeSvc, grpcAddr, grpcCertPath, grpcKeyPath, ca, crl)
} else {
if grpcAddr != "" {
slog.Warn("grpc agent server disabled: one or more required env vars missing",
"CA_KEY_PATH_set", caKeyPath != "",
"CA_CERT_PATH_set", caCertPath != "",
"GRPC_CERT_PATH_set", grpcCertPath != "",
"GRPC_KEY_PATH_set", grpcKeyPath != "")
}
// Even without gRPC, build a Hub-only nodeSvc so the HTTP /v1/nodes
// routes (ListNodes/Connect) work locally. nil CA/tokens means
// agent Enroll would panic if called — acceptable when gRPC is off.
nodeSvc = nodes.NewService(nil, nil, rdb, nodeStore)
nodeSvc.Hub().Start(context.Background())
log.Printf("nodes.Service: gRPC not configured; Hub active for command queueing (HTTP /v1/nodes enabled)")
}
}
// ─── HTTP router ──────────────────────────────────────────────────────────
r := chi.NewRouter()
r.Use(chimw.Logger)
r.Use(chimw.Recoverer)
r.Use(apierr.Middleware)
// CORS:Web 用户中心(pangolin.yanmeiai.com/user)跨域调 /v1/*;原生端不受影响。
r.Use(httpapi.NewCORS())
r.Get("/healthz", func(w http.ResponseWriter, r *http.Request) {
w.Header().Set("Content-Type", "application/json")
w.WriteHeader(http.StatusOK)
_ = json.NewEncoder(w).Encode(map[string]string{"status": "ok"})
})
// Public (no auth): 客户端安装包下载(官网下载按钮直链)。CI
// (scripts/ci/deploy-client.sh)把最新安装包 scp 到 DOWNLOADS_DIR,按平台固定
// 文件名覆盖;目录不存在也不影响启动,只是请求 404(见 DownloadsHandler 注释)。
downloadsHandler := httpapi.NewDownloadsHandler(os.Getenv("DOWNLOADS_DIR"))
r.Get("/downloads/*", downloadsHandler.Serve)
// Public (no auth): 客户端自动更新版本清单。VERSION_MANIFEST 可配置清单路径
// (默认 /etc/pangolin/version.yaml);deploy/single-node/deploy.sh 安装仓库内
// 默认清单,scripts/ci/release-client.sh 在每次 client-v* 发版时改写其
// version/build_number。每次请求都重新读文件,发版脚本改完立即生效,无需重启。
versionHandler := httpapi.NewVersionHandler(os.Getenv("VERSION_MANIFEST"))
r.Get("/version", versionHandler.Serve)
// Optional probe ingest route. sharedProbeStore is reused by the scheduler
// (below) when both are enabled, so they share one Redis-backed store.
var sharedProbeStore *probe.Store
if probeJSON := os.Getenv("PROBE_SECRETS"); probeJSON != "" {
pRDB, err := redisutil.New(redisAddr, os.Getenv("REDIS_PASSWORD"), 0)
if err != nil {
log.Printf("probe: redis connect failed (%v) probe route disabled", err)
} else {
var secretMap map[string]string
if err := json.Unmarshal([]byte(probeJSON), &secretMap); err != nil {
log.Fatalf("probe: invalid PROBE_SECRETS JSON: %v", err)
}
reg := probe.NewMapRegistry(secretMap)
st := probe.NewStore(pRDB)
sharedProbeStore = st
r.Post("/probe/report", probe.NewIngestHandler(reg, st).ServeHTTP)
log.Printf("probe ingest route registered (%d probe(s))", len(secretMap))
}
}
// ─── Scheduler (optional) ─────────────────────────────────────────────────
//
// Set SCHED_ENABLED=true to start the three leader-elected scheduler loops
// (DetectLoop / OrchestrateLoop / CapacityLoop) in this process alongside the
// HTTP and gRPC servers. Rollback: unset SCHED_ENABLED and restart; all other
// server functionality is unaffected.
if os.Getenv("SCHED_ENABLED") == "true" {
schedRDB, err := redisutil.New(redisAddr, os.Getenv("REDIS_PASSWORD"), 0)
if err != nil {
log.Printf("scheduler: redis connect failed (%v) — scheduler disabled", err)
} else {
// Reuse the probe store if the probe route is wired up; otherwise
// create a standalone store on the scheduler Redis client.
ps := sharedProbeStore
if ps == nil {
ps = probe.NewStore(schedRDB)
}
// Prefer real lifecycle (SQL) + provision (#14) wiring when a DB is
// available; fall back to no-op stubs otherwise. Without vendor
// credentials the provision service's CreateNode fails gracefully
// (pending replacements stay pending) while lifecycle reads/transitions
// still operate on real node rows.
var cfg scheduler.Config
if sqlDB != nil {
provStore := provision.NewMySQLStore(sqlDB)
provSvc, perr := provision.NewService(provision.Config{
Store: provStore,
Adapters: providers.NewRegistry(),
})
if perr != nil {
log.Printf("scheduler: provision service init failed (%v) — using stub config", perr)
cfg = scheduler.BuildStubConfig(schedRDB, ps)
} else {
cfg = scheduler.BuildRealConfig(schedRDB, ps, sqlDB,
nodes.NewLoadCache(schedRDB), provSvc, provStore)
log.Printf("scheduler: real lifecycle + provision wired")
}
} else {
cfg = scheduler.BuildStubConfig(schedRDB, ps)
}
sched := scheduler.New(cfg)
// sigCtx cancels on SIGINT/SIGTERM, giving the scheduler up to
// TickTimeout to finish any in-flight tick before process exit.
sigCtx, stopSig := signal.NotifyContext(context.Background(),
os.Interrupt, syscall.SIGTERM)
go func() {
defer stopSig()
if err := sched.Start(sigCtx); err != nil {
slog.Error("scheduler: Start error", "error", err)
}
}()
log.Printf("scheduler started (SCHED_ENABLED=true)")
}
}
// ─── /v1 routes (requires DB) ────────────────────────────────────────────
if sqlDB != nil {
mountV1(r, sqlDB, rdb, nodeSvc)
} else {
log.Printf("DB_DSN not set; /v1 routes disabled")
}
log.Printf("pangolin server listening on %s", *listenAddr)
if err := http.ListenAndServe(*listenAddr, r); err != nil {
log.Fatalf("server error: %v", err)
}
}
// mountV1 wires all /v1 routes onto r.
func mountV1(r chi.Router, sqlDB *sql.DB, rdb *redis.Client, nodeSvc *nodes.Service) {
// ── JWT TokenManager ──────────────────────────────────────────────────────
jwtPrivPath := os.Getenv("JWT_PRIVATE_KEY_PATH")
jwtKID := os.Getenv("JWT_KEY_ID")
var tm *auth.TokenManager
if jwtPrivPath != "" && jwtKID != "" {
tc, err := auth.LoadTokenConfig(jwtPrivPath, jwtKID, parseKeyMap(os.Getenv("JWT_PUBLIC_KEYS")))
if err != nil {
log.Fatalf("JWT keys: %v", err)
}
tm, err = auth.NewTokenManager(rdb, tc)
if err != nil {
log.Fatalf("JWT token manager: %v", err)
}
} else {
log.Printf("JWT not configured — /v1 protected routes will be unavailable")
}
// ── Devices ───────────────────────────────────────────────────────────────
// Constructed before Auth so login/register can register the device.
devicesStore := devices.NewStore(sqlDB)
devicesSvc := devices.NewService(devicesStore, nil) // NoopRevoker for MVP
devicesHandler := devices.NewHandler(devicesSvc)
// Sessions: login sessions bound to (device, refresh JTI) — P2/P3. Shared by
// auth (create/rotate/revoke) and devices (last-login display + force-logout).
sessionStore := sessions.NewStore(sqlDB)
devicesSvc.SetSessionPort(sessionStore)
if tm != nil {
devicesSvc.SetJTIRevoker(tm) // force-logout drops refresh JTIs from Redis
}
// Per-device credential revoke on clear-login: nodes.Service pushes Revoke to
// the holding node(s). Wired when a DB-backed nodeSvc exists.
if nodeSvc != nil {
devicesSvc.SetCredentialRevoker(nodeSvc)
}
// ── Auth ──────────────────────────────────────────────────────────────────
var authHandler *auth.Handler
if tm != nil {
var mailer auth.Mailer
if smtpHost := os.Getenv("SMTP_HOST"); smtpHost != "" {
mailer = auth.NewSMTPMailer(auth.SMTPConfig{
Host: smtpHost,
Port: intEnvDefault("SMTP_PORT", 587),
Username: os.Getenv("SMTP_USERNAME"),
Password: os.Getenv("SMTP_PASSWORD"),
From: getenvDefault("SMTP_FROM", "no-reply@pangolin.app"),
})
} else {
mailer = auth.NewLogMailer(nil) // dev: code printed to log
}
rl := auth.NewRateLimiter(rdb, nil)
authStore := auth.NewSQLStore(sqlDB)
authSvc := auth.NewService(authStore, rdb, rl, tm, mailer, auth.ServiceConfig{}, nil)
authSvc.SetDeviceRegistrar(authDeviceRegistrar{svc: devicesSvc})
authSvc.SetSessionStore(sessionStore)
authHandler = auth.NewHandler(authSvc)
}
// ── User TOTP (web 用户中心 2FA) ─────────────────────────────────────────────
// Only enabled when a valid 32-byte key is configured (USER_TOTP_ENC_KEY,
// raw 32 bytes or 64 hex chars); secrets are encrypted at rest under it.
var totpHandler *auth.TOTPHandler
if tm != nil {
if key := parseTOTPKey(os.Getenv("USER_TOTP_ENC_KEY")); key != nil {
totpHandler = auth.NewTOTPHandler(sqlDB, key, tm, rdb)
} else if os.Getenv("USER_TOTP_ENC_KEY") != "" {
log.Printf("USER_TOTP_ENC_KEY invalid (need 32 bytes or 64 hex chars); user TOTP disabled")
}
}
// ── Codes ────────────────────────────────────────────────────────────────
codesStore := codes.NewStore(sqlDB)
codesSvc := codes.NewService(codesStore, rdb, 5, time.Hour)
redeemHandler := codes.NewRedeemHandler(codesSvc)
webhookHandler := codes.NewWebhookHandler(codesStore, rdb,
os.Getenv("WEBHOOK_SECRET"), 5*time.Minute, 15*time.Minute)
// ── Usage ─────────────────────────────────────────────────────────────────
usageStore := usage.NewStore(sqlDB)
// Ad verifier: real AdMob SSV in production; a放行式 DevVerifier when
// ADS_DEV_MODE=1 (or no AdMob configured) so the placeholder看广告加时 flow
// works end-to-end before the real ad SDK is wired in. Nonce replay
// protection still applies either way.
var adVerifier usage.AdVerifier
if os.Getenv("ADS_DEV_MODE") == "1" {
adVerifier = usage.DevVerifier{}
slog.Warn("ads: DevVerifier enabled (ADS_DEV_MODE=1) — accepts any receipt; not for production")
} else if os.Getenv("ADMOB_SSV") == "1" {
adVerifier = usage.NewAdMobVerifier("", nil, 0, nil)
} else {
// Default (current state): no real ad network yet → placeholder verifier
// so免费版看广告加时 is functional in the field.
adVerifier = usage.DevVerifier{}
slog.Warn("ads: no ad network configured — using DevVerifier placeholder")
}
usageSvc := usage.NewService(usageStore, rdb, adVerifier, time.Hour)
usageHandler := usage.NewUsageHandler(usageSvc)
deviceUsageHandler := usage.NewDeviceUsageHandler(usageSvc)
adsHandler := usage.NewAdsUnlockHandler(usageSvc)
// ── Account / Plans / Notices ─────────────────────────────────────────────
accountAPI := httpapi.NewAccountAPI(sqlDB)
// ── Nodes + Connect ───────────────────────────────────────────────────────
var nodeAPI *httpapi.NodeAPI
var nodeStore nodes.NodeStore
if nodeSvc != nil {
nodeStore = nodeSvc.Store()
publicURL := os.Getenv("PANGOLIN_PUBLIC_URL")
if publicURL == "" {
// 没设 PANGOLIN_PUBLIC_URL → BuildClientConfig 会静默跳过国内分流
// (rule_set 没有公网基址可下载),客户端 split_cn=1 也无效、全量走隧道。
// 这是个隐蔽坑(国内流量绕道出海、变慢),启动期显式告警。
slog.Warn("PANGOLIN_PUBLIC_URL 未设置:国内分流(split_cn)将被静默跳过,客户端全量走隧道。" +
"如需国内直连,设为控制面对外公网基址(如 http://<公网IP>:8080)")
}
nodeAPI = httpapi.NewNodeAPI(nodeStore, nodeSvc.Hub(), nodeSvc.Load(), os.Getenv("NODE_DERIVE_KEY"), publicURL)
}
// 国内分流(#5)的 rule-set 静态服务:GET /v1/rules/{name}.srs(自托管,
// sing-box 直连下载;PANGOLIN_RULES_DIR 默认 /var/lib/pangolin/rules)。
rulesHandler := httpapi.NewRulesHandler(os.Getenv("PANGOLIN_RULES_DIR"))
// ── Subscription URL (web user-center 订阅导入) ───────────────────────────
subAPI := httpapi.NewSubscriptionAPI(sqlDB, nodeStore, os.Getenv("NODE_DERIVE_KEY"), os.Getenv("SUB_BASE"))
// Public subscription endpoint (no JWT): clients fetch their sing-box config
// from the per-user URL. Only available when a node store is configured.
if nodeStore != nil {
r.Get("/sub/{token}", subAPI.ServeSub)
}
// ─── Mount under /v1 ─────────────────────────────────────────────────────
r.Route("/v1", func(v1 chi.Router) {
// Public (no auth): 国内分流 rule-set 下载(sing-box 不带 token)。
v1.Get("/rules/{name}", rulesHandler.Serve)
// Public (no auth): 诊断端点,量节点→目标出海段耗时(白名单限定,防 SSRF),
// 供白盒拆「接入段 vs 出海段」。只返回耗时数字,不传业务数据。
v1.Get("/diag/egress", httpapi.NewDiagHandler().EgressTiming)
// Public (no auth): auth endpoints.
if authHandler != nil {
v1.Post("/auth/code", authHandler.SendCode)
v1.Post("/auth/register", authHandler.Register)
v1.Post("/auth/login", authHandler.Login)
v1.Post("/auth/refresh", authHandler.Refresh)
v1.Post("/auth/logout", authHandler.Logout)
// App→网页免登录换票的公开一端;票据本身即凭证,无需 Bearer(见下方
// 受保护分组里的签票端 /auth/web-ticket)。
v1.Post("/auth/web-ticket/exchange", authHandler.WebTicketExchange)
if totpHandler != nil {
v1.Post("/auth/login/totp", totpHandler.LoginTOTP)
}
}
// Webhook: HMAC-authenticated, no JWT.
v1.Post("/webhook/store/codes", webhookHandler.ServeHTTP)
// Protected: all routes that require a valid Bearer JWT.
if tm != nil {
v1.Group(func(protected chi.Router) {
protected.Use(auth.RequireAuth(tm))
protected.Route("/me", func(me chi.Router) {
me.Get("/", accountAPI.GetMe) // 子路由根,避免与 Route("/me") 冲突致 404
devicesHandler.RegisterRoutes(me)
// Web 用户中心调用 /v1/me/redeemapp 仍用 /v1/redeem。两者同处理器。
me.Post("/redeem", redeemHandler.ServeHTTP)
me.Get("/subscription", subAPI.GetSubscription)
me.Post("/subscription/reset", subAPI.ResetSubscription)
if totpHandler != nil {
me.Post("/totp/setup", totpHandler.Setup)
me.Post("/totp/verify", totpHandler.Verify)
me.Post("/totp/disable", totpHandler.Disable)
}
})
protected.Post("/redeem", redeemHandler.ServeHTTP)
if authHandler != nil {
// 签票端要求已登录(拿当前 JWT 的 uid/uuid);兑换端见上方公开分组。
protected.Post("/auth/web-ticket", authHandler.WebTicket)
}
protected.Get("/usage", usageHandler.ServeHTTP)
protected.Get("/usage/devices", deviceUsageHandler.ServeHTTP)
protected.Post("/ads/unlock", adsHandler.ServeHTTP)
protected.Get("/plans", accountAPI.ListPlans)
protected.Get("/notices", accountAPI.ListNotices)
if nodeAPI != nil {
protected.Get("/nodes", nodeAPI.ListNodes)
protected.Post("/nodes/{id}/connect", nodeAPI.ConnectNode)
protected.Post("/nodes/{id}/disconnect", nodeAPI.DisconnectNode)
}
})
}
})
}
// startAdminIfConfigured launches the admin listener when ADMIN_SECRET_KEY and
// DB_DSN are both present.
func startAdminIfConfigured() {
if os.Getenv("ADMIN_SECRET_KEY") == "" || os.Getenv("DB_DSN") == "" {
log.Printf("admin backend disabled (set DB_DSN and ADMIN_SECRET_KEY to enable)")
return
}
cfg, err := admin.FromEnv()
if err != nil {
log.Fatalf("admin config: %v", err)
}
database, err := db.Open(os.Getenv("DB_DSN"))
if err != nil {
log.Fatalf("admin db: %v", err)
}
rdb, err := redisutil.New(
getenvDefault("REDIS_ADDR", "127.0.0.1:6379"),
os.Getenv("REDIS_PASSWORD"), 0)
if err != nil {
log.Fatalf("admin redis: %v", err)
}
svc := admin.BuildServices(database, rdb, 5, cfg.LoginLockDuration)
handler, err := admin.NewHandler(cfg, database, rdb, svc, log.Default())
if err != nil {
log.Fatalf("admin handler: %v", err)
}
go func() {
log.Printf("admin backend listening on %s (internal only)", cfg.Listen)
if err := http.ListenAndServe(cfg.Listen, handler); err != nil {
log.Fatalf("admin server error: %v", err)
}
}()
}
// startGRPCWithService starts the mTLS gRPC agent server using the shared
// nodes.Service so the HTTP connect handler and the gRPC server share one Hub.
func startGRPCWithService(
svc *nodes.Service,
addr, grpcCertPath, grpcKeyPath string,
ca *mtls.CA,
crl *mtls.CRL,
) {
serverCert, err := tls.LoadX509KeyPair(grpcCertPath, grpcKeyPath)
if err != nil {
log.Fatalf("grpc: load transport TLS cert: %v", err)
}
tlsCfg := mtls.NewServerTLSConfig(ca, crl)
tlsCfg.Certificates = []tls.Certificate{serverCert}
srv := grpc.NewServer(
grpc.Creds(credentials.NewTLS(tlsCfg)),
grpc.ChainUnaryInterceptor(mtls.UnaryServerInterceptor()),
grpc.ChainStreamInterceptor(mtls.StreamServerInterceptor()),
)
agentv1.RegisterAgentServiceServer(srv, svc.Handler())
lis, err := net.Listen("tcp", addr)
if err != nil {
log.Fatalf("grpc: listen %s: %v", addr, err)
}
slog.Info("grpc agent server listening", "addr", addr)
if err := srv.Serve(lis); err != nil {
log.Fatalf("grpc: serve: %v", err)
}
}
// parseKeyMap parses "kid1:/path1,kid2:/path2" into a map. Used for JWT_PUBLIC_KEYS.
func parseKeyMap(raw string) map[string]string {
if raw == "" {
return nil
}
m := map[string]string{}
for _, pair := range strings.Split(raw, ",") {
pair = strings.TrimSpace(pair)
if pair == "" {
continue
}
i := strings.IndexByte(pair, ':')
if i <= 0 || i == len(pair)-1 {
continue
}
m[strings.TrimSpace(pair[:i])] = strings.TrimSpace(pair[i+1:])
}
return m
}
func getenvDefault(key, def string) string {
if v := os.Getenv(key); v != "" {
return v
}
return def
}
// parseTOTPKey returns a 32-byte AES key from USER_TOTP_ENC_KEY, accepting either
// 64 hex chars or a raw 32-byte string. Returns nil when unset/invalid.
func parseTOTPKey(raw string) []byte {
if len(raw) == 64 {
if b, err := hex.DecodeString(raw); err == nil && len(b) == 32 {
return b
}
return nil
}
if len(raw) == 32 {
return []byte(raw)
}
return nil
}
func intEnvDefault(key string, def int) int {
v := os.Getenv(key)
if v == "" {
return def
}
n, err := strconv.Atoi(v)
if err != nil || n == 0 {
return def
}
return n
}
// authDeviceRegistrar adapts devices.Service to auth.DeviceRegistrar, keeping the
// auth and devices packages decoupled. Device registration on login/register is
// best-effort with NO cap enforcement (MaxDevices=0): free-plan reinstall churns
// the device UUID, so a hard cap at login would lock users out. Explicit device
// limiting is a separate future policy with its own UX.
type authDeviceRegistrar struct{ svc *devices.Service }
func (a authDeviceRegistrar) RegisterDevice(ctx context.Context, userID int64, meta auth.DeviceMeta) (int64, error) {
id, _, apiErr := a.svc.RegisterIfAbsent(ctx, devices.RegisterInput{
UserID: userID,
DeviceUUID: meta.DeviceID,
Name: meta.Name,
Platform: meta.Platform,
ClientVersion: meta.ClientVersion,
MaxDevices: 0,
})
if apiErr != nil {
return 0, apiErr
}
return id, nil
}
// CheckDeviceLimit adapts devices.CheckDeviceLimit → auth.DeviceLimit for the login
// over-limit gate. Returns nil (within cap) unless the account is over its plan cap.
func (a authDeviceRegistrar) CheckDeviceLimit(ctx context.Context, userID int64) (*auth.DeviceLimit, error) {
st, apiErr := a.svc.CheckDeviceLimit(ctx, userID)
if apiErr != nil {
return nil, apiErr
}
if st == nil || !st.Over {
return nil, nil
}
briefs := make([]auth.DeviceBrief, 0, len(st.Devices))
for _, d := range st.Devices {
briefs = append(briefs, auth.DeviceBrief{
UUID: d.UUID, Name: d.Name, Platform: d.Platform, LastSeen: d.LastSeen,
})
}
return &auth.DeviceLimit{MaxDevices: st.MaxDevices, Devices: briefs}, nil
}