[优化] 目录调整

This commit is contained in:
ryan
2026-06-06 16:08:40 +08:00
parent 314229b7ee
commit 959b134d67
230 changed files with 468 additions and 296 deletions
@@ -0,0 +1,256 @@
package acme
import (
"crypto"
"crypto/ecdsa"
"crypto/elliptic"
"crypto/rand"
"crypto/rsa"
"crypto/x509"
"encoding/json"
"encoding/pem"
"errors"
"fmt"
"time"
"github.com/go-acme/lego/v4/acme"
"github.com/go-acme/lego/v4/certcrypto"
"github.com/go-acme/lego/v4/certificate"
"github.com/go-acme/lego/v4/challenge/dns01"
"github.com/go-acme/lego/v4/lego"
"github.com/go-acme/lego/v4/providers/dns/cloudflare"
"github.com/go-acme/lego/v4/registration"
)
type AcmeUser struct {
Email string
Registration *registration.Resource
key crypto.PrivateKey
}
func (u *AcmeUser) GetEmail() string {
return u.Email
}
func (u *AcmeUser) GetRegistration() *registration.Resource {
return u.Registration
}
func (u *AcmeUser) GetPrivateKey() crypto.PrivateKey {
return u.key
}
type CertificateResult struct {
CertPEM string
KeyPEM string
NotBefore time.Time
NotAfter time.Time
}
func parsePrivateKey(pemData string) (crypto.PrivateKey, error) {
block, _ := pem.Decode([]byte(pemData))
if block == nil {
return nil, errors.New("failed to parse PEM block containing the key")
}
if key, err := x509.ParsePKCS1PrivateKey(block.Bytes); err == nil {
return key, nil
}
if key, err := x509.ParsePKCS8PrivateKey(block.Bytes); err == nil {
return key, nil
}
if key, err := x509.ParseECPrivateKey(block.Bytes); err == nil {
return key, nil
}
return nil, errors.New("failed to parse private key")
}
func encodePrivateKey(key crypto.PrivateKey) (string, error) {
var pemBlock *pem.Block
switch k := key.(type) {
case *rsa.PrivateKey:
pemBlock = &pem.Block{Type: "RSA PRIVATE KEY", Bytes: x509.MarshalPKCS1PrivateKey(k)}
case *ecdsa.PrivateKey:
b, err := x509.MarshalECPrivateKey(k)
if err != nil {
return "", err
}
pemBlock = &pem.Block{Type: "EC PRIVATE KEY", Bytes: b}
default:
return "", errors.New("unsupported key type")
}
return string(pem.EncodeToMemory(pemBlock)), nil
}
func GetOrCreateLegoClient(acmeEmail, privateKeyPEM, accountURL string, keyAlgorithm string) (*lego.Client, *AcmeUser, string, string, error) {
var privateKey crypto.PrivateKey
var err error
var newPrivateKeyPEM string
var newAccountURL string
if privateKeyPEM == "" {
privateKey, err = ecdsa.GenerateKey(elliptic.P256(), rand.Reader)
if err != nil {
return nil, nil, "", "", err
}
pemStr, err := encodePrivateKey(privateKey)
if err != nil {
return nil, nil, "", "", err
}
newPrivateKeyPEM = pemStr
} else {
privateKey, err = parsePrivateKey(privateKeyPEM)
if err != nil {
return nil, nil, "", "", err
}
}
user := &AcmeUser{
Email: acmeEmail,
key: privateKey,
}
if accountURL != "" {
user.Registration = &registration.Resource{
Body: acme.Account{
Status: "valid",
Contact: []string{"mailto:" + acmeEmail},
},
URI: accountURL,
}
}
config := lego.NewConfig(user)
// Use Let's Encrypt production environment by default
config.CADirURL = lego.LEDirectoryProduction
switch keyAlgorithm {
case "RSA2048":
config.Certificate.KeyType = certcrypto.RSA2048
case "RSA4096":
config.Certificate.KeyType = certcrypto.RSA4096
case "EC256":
config.Certificate.KeyType = certcrypto.EC256
case "EC384":
config.Certificate.KeyType = certcrypto.EC384
default:
config.Certificate.KeyType = certcrypto.RSA2048
}
client, err := lego.NewClient(config)
if err != nil {
return nil, nil, "", "", err
}
if accountURL == "" {
reg, err := client.Registration.Register(registration.RegisterOptions{TermsOfServiceAgreed: true})
if err != nil {
return nil, nil, "", "", err
}
user.Registration = reg
newAccountURL = reg.URI
}
return client, user, newPrivateKeyPEM, newAccountURL, nil
}
func SetupDNSProvider(client *lego.Client, dnsType, dnsAuth string, dns1, dns2 string, disableCNAME, skipDNS bool) error {
var provider challengeProvider
switch dnsType {
case "cloudflare":
var creds map[string]string
if err := json.Unmarshal([]byte(dnsAuth), &creds); err != nil {
return fmt.Errorf("failed to parse cloudflare credentials: %v", err)
}
config := cloudflare.NewDefaultConfig()
config.AuthToken = creds["api_token"]
p, err := cloudflare.NewDNSProviderConfig(config)
if err != nil {
return err
}
provider = p
default:
return fmt.Errorf("unsupported DNS provider: %s", dnsType)
}
var resolvers []string
if dns1 != "" {
resolvers = append(resolvers, dns1+":53")
}
if dns2 != "" {
resolvers = append(resolvers, dns2+":53")
}
var opts []dns01.ChallengeOption
if len(resolvers) > 0 {
opts = append(opts, dns01.AddRecursiveNameservers(resolvers))
}
if disableCNAME {
opts = append(opts, dns01.DisableCompletePropagationRequirement())
}
if skipDNS {
opts = append(opts, dns01.WrapPreCheck(func(domain, fqdn, value string, check dns01.PreCheckFunc) (bool, error) {
time.Sleep(20 * time.Second)
return true, nil
}))
}
return client.Challenge.SetDNS01Provider(provider, opts...)
}
type challengeProvider interface {
Present(domain, token, keyAuth string) error
CleanUp(domain, token, keyAuth string) error
}
func ObtainSSL(
acmeEmail, acmePrivateKeyPEM, acmeURL string,
dnsType, dnsAuth string,
dns1, dns2 string,
disableCNAME, skipDNS bool,
keyAlgorithm string,
domains []string,
) (string, string, *CertificateResult, error) {
client, _, newPrivateKeyPEM, newAccountURL, err := GetOrCreateLegoClient(acmeEmail, acmePrivateKeyPEM, acmeURL, keyAlgorithm)
if err != nil {
return "", "", nil, fmt.Errorf("failed to create ACME client: %w", err)
}
err = SetupDNSProvider(client, dnsType, dnsAuth, dns1, dns2, disableCNAME, skipDNS)
if err != nil {
return newAccountURL, newPrivateKeyPEM, nil, fmt.Errorf("failed to setup DNS provider: %w", err)
}
request := certificate.ObtainRequest{
Domains: domains,
Bundle: true,
}
certificates, err := client.Certificate.Obtain(request)
if err != nil {
return newAccountURL, newPrivateKeyPEM, nil, fmt.Errorf("failed to obtain certificate: %w", err)
}
result := &CertificateResult{
CertPEM: string(certificates.Certificate),
KeyPEM: string(certificates.PrivateKey),
}
// Parse validity dates
certBlock, _ := pem.Decode(certificates.Certificate)
if certBlock != nil {
parsedCert, err := x509.ParseCertificate(certBlock.Bytes)
if err == nil {
result.NotBefore = parsedCert.NotBefore
result.NotAfter = parsedCert.NotAfter
}
}
return newAccountURL, newPrivateKeyPEM, result, nil
}
+221
View File
@@ -0,0 +1,221 @@
package cap
import (
"crypto/hmac"
"crypto/rand"
"crypto/sha256"
"encoding/base64"
"encoding/hex"
"encoding/json"
"errors"
"strconv"
"strings"
"time"
)
const jwtHeaderB64 = "eyJhbGciOiJIUzI1NiIsInR5cCI6IkpXVCJ9"
// ChallengeConfig holds parameters for the PoW challenge
type ChallengeConfig struct {
Count int // Number of puzzles (c)
Size int // Salt length (s)
Difficulty int // Difficulty prefix length (d)
ExpiresMs time.Duration // Challenge TTL
}
// ChallengeResponse is returned to the client
type ChallengeResponse struct {
Challenge struct {
C int `json:"c"`
S int `json:"s"`
D int `json:"d"`
} `json:"challenge"`
Token string `json:"token"`
Expires int64 `json:"expires"` // ms timestamp
}
// ChallengePayload represents the signed JWT payload
type ChallengePayload struct {
Nonce string `json:"n"`
Count int `json:"c"`
Size int `json:"s"`
Difficulty int `json:"d"`
Expires int64 `json:"exp"` // ms timestamp
IssuedAt int64 `json:"iat"` // ms timestamp
Scope string `json:"sk,omitempty"`
}
// RedeemRequest payload sent by client
type RedeemRequest struct {
Token string `json:"token"`
Solutions []int `json:"solutions"`
}
// RedeemResponse returned to client after verification
type RedeemResponse struct {
Success bool `json:"success"`
Token string `json:"token,omitempty"`
Expires int64 `json:"expires,omitempty"`
Error string `json:"error,omitempty"`
}
func b64urlEncode(data []byte) string {
return base64.RawURLEncoding.EncodeToString(data)
}
func b64urlDecode(str string) ([]byte, error) {
return base64.RawURLEncoding.DecodeString(str)
}
func randomHex(byteLen int) string {
bytes := make([]byte, byteLen)
if _, err := rand.Read(bytes); err != nil {
panic(err)
}
return hex.EncodeToString(bytes)
}
func jwtSign(payload []byte, secret []byte) string {
body := b64urlEncode(payload)
sigInput := jwtHeaderB64 + "." + body
mac := hmac.New(sha256.New, secret)
mac.Write([]byte(sigInput))
sig := mac.Sum(nil)
return sigInput + "." + b64urlEncode(sig)
}
func jwtVerify(token string, secret []byte) ([]byte, error) {
parts := strings.Split(token, ".")
if len(parts) != 3 {
return nil, errors.New("invalid token format")
}
if parts[0] != jwtHeaderB64 {
return nil, errors.New("invalid header")
}
sigInput := parts[0] + "." + parts[1]
mac := hmac.New(sha256.New, secret)
mac.Write([]byte(sigInput))
expectedSig := mac.Sum(nil)
actualSig, err := b64urlDecode(parts[2])
if err != nil {
return nil, err
}
if !hmac.Equal(expectedSig, actualSig) {
return nil, errors.New("signature mismatch")
}
payload, err := b64urlDecode(parts[1])
if err != nil {
return nil, err
}
return payload, nil
}
func jwtSigHex(token string) string {
parts := strings.Split(token, ".")
if len(parts) != 3 {
return ""
}
sigBytes, err := b64urlDecode(parts[2])
if err != nil {
return ""
}
return hex.EncodeToString(sigBytes)
}
// GenerateChallenge produces a new challenge and signed token
func GenerateChallenge(secret []byte, conf ChallengeConfig, scope string) (*ChallengeResponse, error) {
if conf.Count <= 0 {
conf.Count = 50
}
if conf.Size <= 0 {
conf.Size = 32
}
if conf.Difficulty <= 0 {
conf.Difficulty = 4
}
if conf.ExpiresMs <= 0 {
conf.ExpiresMs = 10 * time.Minute
}
now := time.Now().UnixNano() / int64(time.Millisecond)
expires := now + int64(conf.ExpiresMs/time.Millisecond)
payload := ChallengePayload{
Nonce: randomHex(25),
Count: conf.Count,
Size: conf.Size,
Difficulty: conf.Difficulty,
Expires: expires,
IssuedAt: now,
Scope: scope,
}
payloadBytes, err := json.Marshal(payload)
if err != nil {
return nil, err
}
token := jwtSign(payloadBytes, secret)
resp := &ChallengeResponse{
Token: token,
Expires: expires,
}
resp.Challenge.C = conf.Count
resp.Challenge.S = conf.Size
resp.Challenge.D = conf.Difficulty
return resp, nil
}
// VerifyChallengeSolutions verifies client submitted solutions
func VerifyChallengeSolutions(token string, solutions []int, secret []byte, expectedScope string) (*ChallengePayload, error) {
payloadBytes, err := jwtVerify(token, secret)
if err != nil {
return nil, errors.New("invalid_token")
}
var payload ChallengePayload
if err := json.Unmarshal(payloadBytes, &payload); err != nil {
return nil, errors.New("invalid_token")
}
if expectedScope != "" && payload.Scope != expectedScope {
return nil, errors.New("scope_mismatch")
}
now := time.Now().UnixNano() / int64(time.Millisecond)
if payload.Expires < now {
return nil, errors.New("expired")
}
if len(solutions) != payload.Count {
return nil, errors.New("invalid_solutions")
}
tokenFnv := fnv1a(token)
for i := 0; i < payload.Count; i++ {
idxStr := strconv.Itoa(i + 1)
saltSeed := fnv1aResume(tokenFnv, idxStr)
targetSeed := fnv1aResume(saltSeed, "d")
salt := prngFromHash(saltSeed, payload.Size)
target := prngFromHash(targetSeed, payload.Difficulty)
hashInput := salt + strconv.Itoa(solutions[i])
hashBytes := sha256.Sum256([]byte(hashInput))
hashHex := hex.EncodeToString(hashBytes[:])
if !strings.HasPrefix(hashHex, target) {
return nil, errors.New("invalid_solution")
}
}
return &payload, nil
}
@@ -0,0 +1,93 @@
package cap
import (
"context"
"crypto/sha256"
"encoding/hex"
"strconv"
"strings"
"testing"
"time"
)
func TestCapFullFlow(t *testing.T) {
secret := []byte("a-very-long-secret-key-at-least-16-bytes")
store := NewMemoryStore(1 * time.Minute)
manager := NewManager(Config{
Secret: secret,
ChallengeCount: 3, // small count for fast test
ChallengeSize: 32,
ChallengeDifficulty: 3, // small difficulty for fast test
ChallengeTTL: 5 * time.Second,
TokenTTL: 10 * time.Second,
}, store)
scope := "test-scope"
resp, err := manager.Generate(scope)
if err != nil {
t.Fatalf("Generate failed: %v", err)
}
if resp.Challenge.C != 3 {
t.Errorf("Expected count 3, got %d", resp.Challenge.C)
}
// Solve the challenge (acting as client)
solutions := make([]int, resp.Challenge.C)
tokenFnv := fnv1a(resp.Token)
for i := 0; i < resp.Challenge.C; i++ {
idxStr := strconv.Itoa(i + 1)
saltSeed := fnv1aResume(tokenFnv, idxStr)
targetSeed := fnv1aResume(saltSeed, "d")
salt := prngFromHash(saltSeed, resp.Challenge.S)
target := prngFromHash(targetSeed, resp.Challenge.D)
// Brute force the PoW solution
var found bool
for nonce := 0; nonce < 1000000; nonce++ {
hashInput := salt + strconv.Itoa(nonce)
hashBytes := sha256.Sum256([]byte(hashInput))
hashHex := hex.EncodeToString(hashBytes[:])
if strings.HasPrefix(hashHex, target) {
solutions[i] = nonce
found = true
break
}
}
if !found {
t.Fatalf("Failed to solve puzzle %d", i)
}
}
// Redeem
ctx := context.Background()
redeemResp, err := manager.Redeem(ctx, resp.Token, solutions, scope)
if err != nil {
t.Fatalf("Redeem failed: %v", err)
}
if !redeemResp.Success {
t.Fatalf("Redeem returned success=false: %s", redeemResp.Error)
}
if redeemResp.Token == "" {
t.Fatalf("Expected token, got empty")
}
// Verify the token
valid, err := manager.VerifyToken(ctx, redeemResp.Token, scope)
if err != nil {
t.Fatalf("VerifyToken failed: %v", err)
}
if !valid {
t.Fatalf("Expected redeem token to be valid")
}
// Verify token is one-time use
validAgain, err := manager.VerifyToken(ctx, redeemResp.Token, scope)
if err != nil {
t.Fatalf("VerifyToken second call failed: %v", err)
}
if validAgain {
t.Fatalf("Expected redeem token to be single-use (invalidated after verification)")
}
}
@@ -0,0 +1,165 @@
package cap
import (
"context"
"crypto/sha256"
"encoding/hex"
"strconv"
"strings"
"time"
)
// Config holds settings for the CAPTCHA manager
type Config struct {
Secret []byte // HMAC signing key
ChallengeCount int // Number of PoW puzzles
ChallengeSize int // Size of the salt string
ChallengeDifficulty int // Length of difficulty target prefix
ChallengeTTL time.Duration // Lifespan of the challenge JWT
TokenTTL time.Duration // Lifespan of the redeem token
}
// Manager orchestrates challenge generation and solution validation
type Manager struct {
conf Config
store Store
}
// NewManager creates a new CAPTCHA Manager
func NewManager(conf Config, store Store) *Manager {
if conf.ChallengeCount <= 0 {
conf.ChallengeCount = 50
}
if conf.ChallengeSize <= 0 {
conf.ChallengeSize = 32
}
if conf.ChallengeDifficulty <= 0 {
conf.ChallengeDifficulty = 4
}
if conf.ChallengeTTL <= 0 {
conf.ChallengeTTL = 10 * time.Minute
}
if conf.TokenTTL <= 0 {
conf.TokenTTL = 20 * time.Minute
}
return &Manager{
conf: conf,
store: store,
}
}
// Generate creates a challenge response
func (m *Manager) Generate(scope string) (*ChallengeResponse, error) {
c := ChallengeConfig{
Count: m.conf.ChallengeCount,
Size: m.conf.ChallengeSize,
Difficulty: m.conf.ChallengeDifficulty,
ExpiresMs: m.conf.ChallengeTTL,
}
return GenerateChallenge(m.conf.Secret, c, scope)
}
// Redeem verifies PoW solutions and returns a one-time redeem token
func (m *Manager) Redeem(ctx context.Context, token string, solutions []int, scope string) (*RedeemResponse, error) {
sigHex := jwtSigHex(token)
if sigHex == "" {
return &RedeemResponse{Success: false, Error: "invalid_token"}, nil
}
// Replay prevention: check if this JWT signature has already been used
nonceKey := "cap:nonce:" + sigHex
_, exists, err := m.store.Get(ctx, nonceKey)
if err != nil {
return &RedeemResponse{Success: false, Error: "nonce_store_error"}, err
}
if exists {
return &RedeemResponse{Success: false, Error: "already_redeemed"}, nil
}
payload, err := VerifyChallengeSolutions(token, solutions, m.conf.Secret, scope)
if err != nil {
return &RedeemResponse{Success: false, Error: err.Error()}, nil
}
// Verification succeeded. Consume the nonce.
now := time.Now().UnixNano() / int64(time.Millisecond)
ttlMs := time.Duration(payload.Expires-now) * time.Millisecond
if ttlMs < time.Second {
ttlMs = time.Second
}
if err := m.store.Set(ctx, nonceKey, "1", ttlMs); err != nil {
return &RedeemResponse{Success: false, Error: "nonce_store_error"}, err
}
// Generate a redeem token formatted as "id:verToken"
id := randomHex(8)
verToken := randomHex(15)
verHashBytes := sha256.Sum256([]byte(verToken))
verHashHex := hex.EncodeToString(verHashBytes[:])
tokenKey := "cap:token:" + id + ":" + verHashHex
tokenExpires := time.Now().Add(m.conf.TokenTTL)
// Value stored is "expiresNano|scope"
storeVal := strconv.FormatInt(tokenExpires.UnixNano(), 10) + "|" + scope
if err := m.store.Set(ctx, tokenKey, storeVal, m.conf.TokenTTL); err != nil {
return &RedeemResponse{Success: false, Error: "token_store_error"}, err
}
return &RedeemResponse{
Success: true,
Token: id + ":" + verToken,
Expires: tokenExpires.UnixNano() / int64(time.Millisecond),
}, nil
}
// VerifyToken validates and consumes the redeem token (single-use)
func (m *Manager) VerifyToken(ctx context.Context, token string, expectedScope string) (bool, error) {
if token == "" {
return false, nil
}
parts := strings.Split(token, ":")
if len(parts) != 2 {
return false, nil
}
id := parts[0]
verToken := parts[1]
verHashBytes := sha256.Sum256([]byte(verToken))
verHashHex := hex.EncodeToString(verHashBytes[:])
tokenKey := "cap:token:" + id + ":" + verHashHex
val, exists, err := m.store.Get(ctx, tokenKey)
if err != nil {
return false, err
}
if !exists {
return false, nil
}
// Single-use: consume/delete the token immediately
_ = m.store.Delete(ctx, tokenKey)
valParts := strings.Split(val, "|")
if len(valParts) != 2 {
return false, nil
}
expNano, err := strconv.ParseInt(valParts[0], 10, 64)
if err != nil {
return false, nil
}
tokenScope := valParts[1]
if expectedScope != "" && tokenScope != expectedScope {
return false, nil
}
if time.Now().UnixNano() > expNano {
return false, nil // Expired
}
return true, nil
}
@@ -0,0 +1,40 @@
package cap
import (
"net/http"
"github.com/gin-gonic/gin"
)
// VerifyMiddleware returns a Gin middleware that checks and consumes the X-Cap-Token header.
// enabledFunc is an optional callback allowing dynamic check of whether captcha protection is turned on.
func (m *Manager) VerifyMiddleware(scope string, enabledFunc func() bool) gin.HandlerFunc {
return func(c *gin.Context) {
if enabledFunc != nil && !enabledFunc() {
c.Next()
return
}
token := c.GetHeader("X-Cap-Token")
if token == "" {
c.JSON(http.StatusUnauthorized, gin.H{
"success": false,
"error": "验证码验证失败,缺少验证码凭证",
})
c.Abort()
return
}
valid, err := m.VerifyToken(c.Request.Context(), token, scope)
if err != nil || !valid {
c.JSON(http.StatusUnauthorized, gin.H{
"success": false,
"error": "验证码校验失败或已过期,请重试",
})
c.Abort()
return
}
c.Next()
}
}
@@ -0,0 +1,45 @@
package cap
import (
"fmt"
"strings"
)
// fnv1a returns the 32-bit FNV-1a hash of a string
func fnv1a(str string) uint32 {
var hash uint32 = 2166136261
for i := 0; i < len(str); i++ {
hash ^= uint32(str[i])
hash += (hash << 1) + (hash << 4) + (hash << 7) + (hash << 8) + (hash << 24)
}
return hash
}
// fnv1aResume resumes FNV-1a hashing from a given state
func fnv1aResume(state uint32, str string) uint32 {
h := state
for i := 0; i < len(str); i++ {
h ^= uint32(str[i])
h += (h << 1) + (h << 4) + (h << 7) + (h << 8) + (h << 24)
}
return h
}
// prng generates a hex string of specified length using a seed
func prng(seed string, length int) string {
return prngFromHash(fnv1a(seed), length)
}
// prngFromHash generates a hex string of specified length using an initial hash state
func prngFromHash(initialHash uint32, length int) string {
state := initialHash
var result strings.Builder
for result.Len() < length {
state ^= state << 13
state ^= state >> 17
state ^= state << 5
hexStr := fmt.Sprintf("%08x", state)
result.WriteString(hexStr)
}
return result.String()[:length]
}
@@ -0,0 +1,90 @@
package cap
import (
"context"
"sync"
"time"
)
// Store defines the storage interface for challenge nonces and verification tokens
type Store interface {
Get(ctx context.Context, key string) (string, bool, error)
Set(ctx context.Context, key string, val string, ttl time.Duration) error
Delete(ctx context.Context, key string) error
}
type memoryItem struct {
value string
expiresAt time.Time
}
// MemoryStore is a thread-safe in-memory implementation of Store
type MemoryStore struct {
items map[string]memoryItem
mu sync.RWMutex
}
// NewMemoryStore creates and initializes a new MemoryStore
func NewMemoryStore(cleanupInterval time.Duration) *MemoryStore {
store := &MemoryStore{
items: make(map[string]memoryItem),
}
if cleanupInterval > 0 {
go store.startCleanupLoop(cleanupInterval)
}
return store
}
func (s *MemoryStore) Get(ctx context.Context, key string) (string, bool, error) {
s.mu.RLock()
item, found := s.items[key]
s.mu.RUnlock()
if !found {
return "", false, nil
}
if time.Now().After(item.expiresAt) {
s.mu.Lock()
delete(s.items, key)
s.mu.Unlock()
return "", false, nil
}
return item.value, true, nil
}
func (s *MemoryStore) Set(ctx context.Context, key string, val string, ttl time.Duration) error {
s.mu.Lock()
defer s.mu.Unlock()
s.items[key] = memoryItem{
value: val,
expiresAt: time.Now().Add(ttl),
}
return nil
}
func (s *MemoryStore) Delete(ctx context.Context, key string) error {
s.mu.Lock()
defer s.mu.Unlock()
delete(s.items, key)
return nil
}
func (s *MemoryStore) startCleanupLoop(interval time.Duration) {
ticker := time.NewTicker(interval)
for range ticker.C {
s.cleanupExpired()
}
}
func (s *MemoryStore) cleanupExpired() {
now := time.Now()
s.mu.Lock()
defer s.mu.Unlock()
for k, v := range s.items {
if now.After(v.expiresAt) {
delete(s.items, k)
}
}
}
@@ -0,0 +1,40 @@
package embedfs
import (
"embed"
"io/fs"
"net/http"
"strings"
"github.com/gin-contrib/static"
)
// Credit: https://github.com/gin-contrib/static/issues/19
type fileSystem struct {
http.FileSystem
}
func (e fileSystem) Exists(prefix string, path string) bool {
cleanPath := strings.TrimPrefix(path, prefix)
cleanPath = strings.TrimPrefix(cleanPath, "/")
if cleanPath == "" {
return false
}
_, err := e.Open(cleanPath)
if err != nil {
return false
}
return true
}
func EmbedFolder(fsEmbed embed.FS, targetPath string) static.ServeFileSystem {
efs, err := fs.Sub(fsEmbed, targetPath)
if err != nil {
panic(err)
}
return fileSystem{
FileSystem: http.FS(efs),
}
}
@@ -0,0 +1,74 @@
package mail
import (
"crypto/tls"
"encoding/base64"
"fmt"
"net/smtp"
"strings"
)
// SMTPConfig holds all the configuration parameters required to send an email.
type SMTPConfig struct {
Server string
Port int
Account string
Token string
SystemName string
}
// SendEmail sends an HTML email to the receiver using the provided SMTP configuration.
func SendEmail(config SMTPConfig, subject string, receiver string, content string) error {
encodedSubject := fmt.Sprintf("=?UTF-8?B?%s?=", base64.StdEncoding.EncodeToString([]byte(subject)))
mail := []byte(fmt.Sprintf("To: %s\r\n"+
"From: %s<%s>\r\n"+
"Subject: %s\r\n"+
"Content-Type: text/html; charset=UTF-8\r\n\r\n%s\r\n",
receiver, config.SystemName, config.Account, encodedSubject, content))
auth := smtp.PlainAuth("", config.Account, config.Token, config.Server)
addr := fmt.Sprintf("%s:%d", config.Server, config.Port)
to := strings.Split(receiver, ";")
var err error
if config.Port == 465 {
tlsConfig := &tls.Config{
InsecureSkipVerify: true,
ServerName: config.Server,
}
conn, err := tls.Dial("tcp", fmt.Sprintf("%s:%d", config.Server, config.Port), tlsConfig)
if err != nil {
return err
}
client, err := smtp.NewClient(conn, config.Server)
if err != nil {
return err
}
defer client.Close()
if err = client.Auth(auth); err != nil {
return err
}
if err = client.Mail(config.Account); err != nil {
return err
}
receiverEmails := strings.Split(receiver, ";")
for _, r := range receiverEmails {
if err = client.Rcpt(r); err != nil {
return err
}
}
w, err := client.Data()
if err != nil {
return err
}
_, err = w.Write(mail)
if err != nil {
return err
}
err = w.Close()
if err != nil {
return err
}
} else {
err = smtp.SendMail(addr, auth, config.Account, to, mail)
}
return err
}
@@ -0,0 +1,67 @@
package ratelimit
import (
"sync"
"time"
)
type InMemoryRateLimiter struct {
store map[string]*[]int64
mutex sync.Mutex
expirationDuration time.Duration
}
func (l *InMemoryRateLimiter) Init(expirationDuration time.Duration) {
if l.store == nil {
l.mutex.Lock()
if l.store == nil {
l.store = make(map[string]*[]int64)
l.expirationDuration = expirationDuration
if expirationDuration > 0 {
go l.clearExpiredItems()
}
}
l.mutex.Unlock()
}
}
func (l *InMemoryRateLimiter) clearExpiredItems() {
for {
time.Sleep(l.expirationDuration)
l.mutex.Lock()
now := time.Now().Unix()
for key := range l.store {
queue := l.store[key]
size := len(*queue)
if size == 0 || now-(*queue)[size-1] > int64(l.expirationDuration.Seconds()) {
delete(l.store, key)
}
}
l.mutex.Unlock()
}
}
// Request parameter duration's unit is seconds
func (l *InMemoryRateLimiter) Request(key string, maxRequestNum int, duration int64) bool {
l.mutex.Lock()
defer l.mutex.Unlock()
// [old <-- new]
queue, ok := l.store[key]
now := time.Now().Unix()
if ok {
if len(*queue) < maxRequestNum {
*queue = append(*queue, now)
return true
}
if now-(*queue)[0] >= duration {
*queue = (*queue)[1:]
*queue = append(*queue, now)
return true
}
return false
}
s := make([]int64, 0, maxRequestNum)
l.store[key] = &s
*(l.store[key]) = append(*(l.store[key]), now)
return true
}
@@ -0,0 +1,14 @@
package security
import "golang.org/x/crypto/bcrypt"
func Password2Hash(password string) (string, error) {
passwordBytes := []byte(password)
hashedPassword, err := bcrypt.GenerateFromPassword(passwordBytes, bcrypt.DefaultCost)
return string(hashedPassword), err
}
func ValidatePasswordAndHash(password string, hash string) bool {
err := bcrypt.CompareHashAndPassword([]byte(hash), []byte(password))
return err == nil
}
@@ -0,0 +1,37 @@
package security
import "crypto/rand"
func GenerateRandomString(length int) string {
if length <= 0 {
return ""
}
const charset = "0123456789ABCDEFGHIJKLMNOPQRSTUVWXYZabcdefghijklmnopqrstuvwxyz"
const n = byte(len(charset))
const threshold = byte(256 - (256 % len(charset)))
out := make([]byte, 0, length)
buf := make([]byte, length)
for len(out) < length {
if _, err := rand.Read(buf); err != nil {
return ""
}
for _, b := range buf {
if b < threshold {
out = append(out, charset[int(b%n)])
if len(out) == length {
break
}
}
}
}
return string(out)
}
func GeneratePassword() string {
return GenerateRandomString(12)
}
func GenerateToken() string {
return GenerateRandomString(22)
}
@@ -0,0 +1,78 @@
package security
import (
"strings"
"sync"
"time"
"github.com/google/uuid"
)
type verificationValue struct {
code string
time time.Time
}
const (
EmailVerificationPurpose = "v"
PasswordResetPurpose = "r"
)
var verificationMutex sync.Mutex
var verificationMap map[string]verificationValue
var verificationMapMaxSize = 10
var VerificationValidMinutes = 10
func GenerateVerificationCode(length int) string {
code := uuid.New().String()
code = strings.Replace(code, "-", "", -1)
if length == 0 {
return code
}
return code[:length]
}
func RegisterVerificationCodeWithKey(key string, code string, purpose string) {
verificationMutex.Lock()
defer verificationMutex.Unlock()
verificationMap[purpose+key] = verificationValue{
code: code,
time: time.Now(),
}
if len(verificationMap) > verificationMapMaxSize {
removeExpiredPairs()
}
}
func VerifyCodeWithKey(key string, code string, purpose string) bool {
verificationMutex.Lock()
defer verificationMutex.Unlock()
value, okay := verificationMap[purpose+key]
now := time.Now()
if !okay || int(now.Sub(value.time).Seconds()) >= VerificationValidMinutes*60 {
return false
}
return code == value.code
}
func DeleteKey(key string, purpose string) {
verificationMutex.Lock()
defer verificationMutex.Unlock()
delete(verificationMap, purpose+key)
}
// no lock inside, so the caller must lock the verificationMap before calling!
func removeExpiredPairs() {
now := time.Now()
for key := range verificationMap {
if int(now.Sub(verificationMap[key].time).Seconds()) >= VerificationValidMinutes*60 {
delete(verificationMap, key)
}
}
}
func init() {
verificationMutex.Lock()
defer verificationMutex.Unlock()
verificationMap = make(map[string]verificationValue)
}
@@ -0,0 +1,373 @@
package uptimekuma
import (
"context"
"encoding/json"
"fmt"
"io"
"log/slog"
"net/http"
"strconv"
"strings"
"sync"
"time"
)
type UptimeKumaMonitor struct {
ID int `json:"id"`
Name string `json:"name"`
Url string `json:"url"`
Type string `json:"type"`
Interval int `json:"interval"`
MaxRetries int `json:"maxretries"`
RetryInterval int `json:"retryInterval"`
Timeout int `json:"timeout"`
Tags []UptimeKumaTag `json:"tags"`
}
type UptimeKumaTag struct {
ID int `json:"tag_id"`
Name string `json:"name"`
Color string `json:"color"`
}
type UptimeKumaTagItem struct {
ID int `json:"id"`
Name string `json:"name"`
Color string `json:"color"`
}
type SocketIOClient struct {
baseURL string
httpClient *http.Client
sid string
ackMutex sync.Mutex
ackID int
ackChanMap map[int]chan string
doneChan chan struct{}
closeOnce sync.Once
monitorListMutex sync.RWMutex
monitorList map[string]UptimeKumaMonitor
monitorListChan chan struct{}
monitorListOnce sync.Once
ctx context.Context
cancel context.CancelFunc
err error
}
func NewSocketIOClient(baseURL string) *SocketIOClient {
ctx, cancel := context.WithCancel(context.Background())
return &SocketIOClient{
baseURL: strings.TrimSuffix(baseURL, "/"),
httpClient: &http.Client{
Timeout: 60 * time.Second,
},
ackChanMap: make(map[int]chan string),
doneChan: make(chan struct{}),
monitorListChan: make(chan struct{}),
monitorList: make(map[string]UptimeKumaMonitor),
ctx: ctx,
cancel: cancel,
}
}
func (c *SocketIOClient) Connect() error {
slog.Debug("Uptime Kuma client starting handshake", "baseURL", c.baseURL)
// 1. Handshake
u := fmt.Sprintf("%s/socket.io/?EIO=4&transport=polling", c.baseURL)
reqHandshake, err := http.NewRequestWithContext(c.ctx, "GET", u, nil)
if err != nil {
return fmt.Errorf("create handshake request failed: %w", err)
}
resp, err := c.httpClient.Do(reqHandshake)
if err != nil {
slog.Error("Uptime Kuma handshake connection failed", "url", u, "error", err)
return fmt.Errorf("handshake request failed: %w", err)
}
defer resp.Body.Close()
bs, err := io.ReadAll(resp.Body)
if err != nil {
slog.Error("Failed to read Uptime Kuma handshake response body", "error", err)
return fmt.Errorf("read handshake body failed: %w", err)
}
bodyStr := string(bs)
slog.Debug("Received handshake response from Uptime Kuma", "body", bodyStr)
if len(bodyStr) == 0 || bodyStr[0] != '0' {
return fmt.Errorf("invalid handshake response format: %s", bodyStr)
}
var hs struct {
Sid string `json:"sid"`
}
if err := json.Unmarshal([]byte(bodyStr[1:]), &hs); err != nil {
return fmt.Errorf("unmarshal handshake sid failed: %w", err)
}
c.sid = hs.Sid
slog.Debug("Uptime Kuma handshake success", "sid", c.sid)
// 2. Namespace Connect
slog.Debug("Sending namespace connect request to Uptime Kuma", "sid", c.sid)
connectURL := fmt.Sprintf("%s/socket.io/?EIO=4&transport=polling&sid=%s", c.baseURL, c.sid)
req, err := http.NewRequestWithContext(c.ctx, "POST", connectURL, strings.NewReader("40"))
if err != nil {
return fmt.Errorf("create connect request failed: %w", err)
}
req.Header.Set("Content-Type", "text/plain;charset=UTF-8")
respConnect, err := c.httpClient.Do(req)
if err != nil {
slog.Error("Uptime Kuma namespace connect request failed", "sid", c.sid, "error", err)
return fmt.Errorf("namespace connect failed: %w", err)
}
respConnect.Body.Close()
slog.Debug("Namespace connected successfully to Uptime Kuma", "sid", c.sid)
// 3. Start Polling Loop
go c.pollLoop()
return nil
}
func (c *SocketIOClient) pollLoop() {
slog.Debug("Uptime Kuma polling loop started", "sid", c.sid)
defer c.Close()
for {
select {
case <-c.doneChan:
slog.Debug("Uptime Kuma polling loop stopped (doneChan closed)", "sid", c.sid)
return
default:
}
u := fmt.Sprintf("%s/socket.io/?EIO=4&transport=polling&sid=%s", c.baseURL, c.sid)
reqPoll, err := http.NewRequestWithContext(c.ctx, "GET", u, nil)
if err != nil {
slog.Error("Failed to create Uptime Kuma polling request", "sid", c.sid, "error", err)
c.err = err
return
}
resp, err := c.httpClient.Do(reqPoll)
if err != nil {
slog.Error("Uptime Kuma polling request failed", "sid", c.sid, "error", err)
c.err = err
return
}
bs, err := io.ReadAll(resp.Body)
resp.Body.Close()
if err != nil {
slog.Error("Failed to read Uptime Kuma polling body", "sid", c.sid, "error", err)
c.err = err
return
}
bodyStr := string(bs)
if len(bodyStr) == 0 {
continue
}
slog.Debug("Received polling payload from Uptime Kuma", "length", len(bodyStr))
packets := strings.Split(bodyStr, "\x1e")
for _, pkt := range packets {
if len(pkt) == 0 {
continue
}
engineIOType := pkt[0]
payload := pkt[1:]
slog.Debug("Parsing engine.io packet", "type", string(engineIOType), "payload_len", len(payload))
switch engineIOType {
case '2': // Ping
slog.Debug("Received engine.io ping, responding with pong", "sid", c.sid)
c.sendPong()
case '4': // Message
if len(payload) == 0 {
continue
}
socketIOType := payload[0]
socketIOPayload := payload[1:]
slog.Debug("Parsing socket.io packet", "type", string(socketIOType), "payload", socketIOPayload)
switch socketIOType {
case '2': // Event
c.handleEvent(socketIOPayload)
case '3': // Ack
c.handleAck(socketIOPayload)
}
}
}
}
}
func (c *SocketIOClient) sendPong() {
u := fmt.Sprintf("%s/socket.io/?EIO=4&transport=polling&sid=%s", c.baseURL, c.sid)
req, err := http.NewRequestWithContext(c.ctx, "POST", u, strings.NewReader("3"))
if err != nil {
return
}
req.Header.Set("Content-Type", "text/plain;charset=UTF-8")
resp, err := c.httpClient.Do(req)
if err == nil {
resp.Body.Close()
}
}
func (c *SocketIOClient) handleEvent(payload string) {
var arr []json.RawMessage
if err := json.Unmarshal([]byte(payload), &arr); err != nil || len(arr) < 2 {
return
}
var eventName string
if err := json.Unmarshal(arr[0], &eventName); err != nil {
return
}
if eventName == "monitorList" {
var list map[string]UptimeKumaMonitor
if err := json.Unmarshal(arr[1], &list); err == nil {
c.monitorListMutex.Lock()
c.monitorList = list
c.monitorListMutex.Unlock()
c.monitorListOnce.Do(func() {
close(c.monitorListChan)
})
}
}
}
func (c *SocketIOClient) handleAck(payload string) {
idx := strings.IndexByte(payload, '[')
if idx == -1 {
return
}
ackIDStr := payload[:idx]
ackID, err := strconv.Atoi(ackIDStr)
if err != nil {
return
}
c.ackMutex.Lock()
ch, ok := c.ackChanMap[ackID]
if ok {
delete(c.ackChanMap, ackID)
c.ackMutex.Unlock()
select {
case ch <- payload[idx:]:
default:
}
} else {
c.ackMutex.Unlock()
}
}
func (c *SocketIOClient) Emit(event string, args ...any) (string, error) {
c.ackMutex.Lock()
id := c.ackID
c.ackID++
ch := make(chan string, 1)
c.ackChanMap[id] = ch
c.ackMutex.Unlock()
payloadArr := []any{event}
payloadArr = append(payloadArr, args...)
bs, err := json.Marshal(payloadArr)
if err != nil {
c.ackMutex.Lock()
delete(c.ackChanMap, id)
c.ackMutex.Unlock()
slog.Error("Failed to marshal event payload", "event", event, "error", err)
return "", err
}
body := fmt.Sprintf("42%d%s", id, string(bs))
slog.Debug("Emitting Socket.IO event", "event", event, "ackID", id, "payload", string(bs))
u := fmt.Sprintf("%s/socket.io/?EIO=4&transport=polling&sid=%s", c.baseURL, c.sid)
req, err := http.NewRequestWithContext(c.ctx, "POST", u, strings.NewReader(body))
if err != nil {
c.ackMutex.Lock()
delete(c.ackChanMap, id)
c.ackMutex.Unlock()
return "", err
}
req.Header.Set("Content-Type", "text/plain;charset=UTF-8")
resp, err := c.httpClient.Do(req)
if err != nil {
c.ackMutex.Lock()
delete(c.ackChanMap, id)
c.ackMutex.Unlock()
slog.Error("Failed to send Emit request", "event", event, "ackID", id, "error", err)
return "", err
}
resp.Body.Close()
select {
case result := <-ch:
slog.Debug("Received Ack for event", "event", event, "ackID", id, "response", result)
return result, nil
case <-time.After(10 * time.Second):
c.ackMutex.Lock()
delete(c.ackChanMap, id)
c.ackMutex.Unlock()
slog.Error("Timeout waiting for event Ack", "event", event, "ackID", id)
return "", fmt.Errorf("timeout waiting for ack for event: %s", event)
case <-c.doneChan:
c.ackMutex.Lock()
delete(c.ackChanMap, id)
c.ackMutex.Unlock()
slog.Error("Client closed while waiting for event Ack", "event", event, "ackID", id)
return "", fmt.Errorf("client closed while waiting for event ack: %s", event)
}
}
func (c *SocketIOClient) Close() {
c.closeOnce.Do(func() {
c.cancel()
close(c.doneChan)
})
}
func (c *SocketIOClient) GetMonitorListChan() <-chan struct{} {
return c.monitorListChan
}
func (c *SocketIOClient) GetMonitorList() map[string]UptimeKumaMonitor {
c.monitorListMutex.RLock()
defer c.monitorListMutex.RUnlock()
// Return a copy to prevent concurrent map read/write access
m := make(map[string]UptimeKumaMonitor, len(c.monitorList))
for k, v := range c.monitorList {
m[k] = v
}
return m
}
func ParseAckResponse(response string, target any) error {
var arr []json.RawMessage
if err := json.Unmarshal([]byte(response), &arr); err != nil || len(arr) == 0 {
return fmt.Errorf("invalid ack response format: %s", response)
}
var status struct {
Ok bool `json:"ok"`
Msg string `json:"msg"`
}
if err := json.Unmarshal(arr[0], &status); err == nil {
if !status.Ok {
errMsg := status.Msg
if errMsg == "" {
errMsg = "unknown error from Uptime Kuma"
}
return fmt.Errorf("Uptime Kuma error response: %s", errMsg)
}
}
if target != nil {
return json.Unmarshal(arr[0], target)
}
return nil
}
@@ -0,0 +1,9 @@
package validation
import "github.com/go-playground/validator/v10"
var Validate *validator.Validate
func init() {
Validate = validator.New()
}