用户登录和接口鉴权
This commit is contained in:
@@ -10,7 +10,9 @@ import (
|
||||
"time"
|
||||
|
||||
"git.huangwc.com/pig/pig-farm-controller/internal/config"
|
||||
"git.huangwc.com/pig/pig-farm-controller/internal/controller/user"
|
||||
"git.huangwc.com/pig/pig-farm-controller/internal/logs"
|
||||
"git.huangwc.com/pig/pig-farm-controller/internal/storage/repository"
|
||||
"github.com/gin-gonic/gin"
|
||||
)
|
||||
|
||||
@@ -26,13 +28,16 @@ type API struct {
|
||||
// config 应用配置
|
||||
config *config.Config
|
||||
|
||||
// userController 用户控制器
|
||||
userController *user.Controller
|
||||
|
||||
// logger 日志记录器
|
||||
logger *logs.Logger
|
||||
}
|
||||
|
||||
// NewAPI 创建并返回一个新的API实例
|
||||
// 初始化Gin引擎和相关配置
|
||||
func NewAPI(cfg *config.Config) *API {
|
||||
func NewAPI(cfg *config.Config, userRepo repository.UserRepo) *API {
|
||||
// 设置Gin为发布模式
|
||||
gin.SetMode(gin.ReleaseMode)
|
||||
|
||||
@@ -56,10 +61,14 @@ func NewAPI(cfg *config.Config) *API {
|
||||
|
||||
engine.Use(gin.Recovery())
|
||||
|
||||
// 创建用户控制器
|
||||
userController := user.NewController(userRepo)
|
||||
|
||||
return &API{
|
||||
engine: engine,
|
||||
config: cfg,
|
||||
logger: logs.NewLogger(),
|
||||
engine: engine,
|
||||
config: cfg,
|
||||
userController: userController,
|
||||
logger: logs.NewLogger(),
|
||||
}
|
||||
}
|
||||
|
||||
@@ -79,11 +88,11 @@ func (a *API) Start() error {
|
||||
}
|
||||
|
||||
// 启动HTTP服务器
|
||||
a.logger.Info(fmt.Sprintf("Starting HTTP server on %s:%d", a.config.Server.Host, a.config.Server.Port))
|
||||
a.logger.Info(fmt.Sprintf("正在启动HTTP服务器 %s:%d", a.config.Server.Host, a.config.Server.Port))
|
||||
|
||||
go func() {
|
||||
if err := a.server.ListenAndServe(); err != nil && err != http.ErrServerClosed {
|
||||
a.logger.Error(fmt.Sprintf("HTTP server startup failed: %v", err))
|
||||
a.logger.Error(fmt.Sprintf("HTTP服务器启动失败: %v", err))
|
||||
}
|
||||
}()
|
||||
|
||||
@@ -92,7 +101,7 @@ func (a *API) Start() error {
|
||||
|
||||
// Stop 停止HTTP服务器
|
||||
func (a *API) Stop() error {
|
||||
a.logger.Info("Stopping HTTP server")
|
||||
a.logger.Info("正在停止HTTP服务器")
|
||||
|
||||
// 创建一个5秒的超时上下文
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
|
||||
@@ -100,11 +109,11 @@ func (a *API) Stop() error {
|
||||
|
||||
// 优雅地关闭服务器
|
||||
if err := a.server.Shutdown(ctx); err != nil {
|
||||
a.logger.Error(fmt.Sprintf("HTTP server shutdown error: %v", err))
|
||||
a.logger.Error(fmt.Sprintf("HTTP服务器关闭错误: %v", err))
|
||||
return err
|
||||
}
|
||||
|
||||
a.logger.Info("HTTP server stopped")
|
||||
a.logger.Info("HTTP服务器已停止")
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -113,6 +122,13 @@ func (a *API) setupRoutes() {
|
||||
// 基础路由示例
|
||||
a.engine.GET("/health", a.healthHandler)
|
||||
|
||||
// 用户相关路由
|
||||
userGroup := a.engine.Group("/api/v1/user")
|
||||
{
|
||||
userGroup.POST("/register", a.userController.Register)
|
||||
userGroup.POST("/login", a.userController.Login)
|
||||
}
|
||||
|
||||
// TODO: 添加更多路由
|
||||
}
|
||||
|
||||
@@ -127,6 +143,6 @@ func (a *API) setupRoutes() {
|
||||
func (a *API) healthHandler(c *gin.Context) {
|
||||
c.JSON(http.StatusOK, gin.H{
|
||||
"status": "ok",
|
||||
"message": "Pig Farm Controller API is running",
|
||||
"message": "猪场控制器API正在运行",
|
||||
})
|
||||
}
|
||||
|
||||
148
internal/api/middleware/auth.go
Normal file
148
internal/api/middleware/auth.go
Normal file
@@ -0,0 +1,148 @@
|
||||
// Package middleware 提供HTTP中间件功能
|
||||
// 包含鉴权、日志、恢复等中间件实现
|
||||
package middleware
|
||||
|
||||
import (
|
||||
"net/http"
|
||||
"os"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"git.huangwc.com/pig/pig-farm-controller/internal/logs"
|
||||
"git.huangwc.com/pig/pig-farm-controller/internal/storage/repository"
|
||||
"github.com/gin-gonic/gin"
|
||||
"github.com/golang-jwt/jwt/v5"
|
||||
"gorm.io/gorm"
|
||||
)
|
||||
|
||||
// AuthMiddleware 鉴权中间件结构
|
||||
type AuthMiddleware struct {
|
||||
userRepo repository.UserRepo
|
||||
logger *logs.Logger
|
||||
}
|
||||
|
||||
// AuthUser 用于在上下文中存储的用户信息
|
||||
type AuthUser struct {
|
||||
ID uint `json:"id"`
|
||||
Username string `json:"username"`
|
||||
}
|
||||
|
||||
// JWTClaims 自定义JWT声明
|
||||
type JWTClaims struct {
|
||||
UserID uint `json:"user_id"`
|
||||
Username string `json:"username"`
|
||||
jwt.RegisteredClaims
|
||||
}
|
||||
|
||||
// NewAuthMiddleware 创建鉴权中间件实例
|
||||
func NewAuthMiddleware(userRepo repository.UserRepo) *AuthMiddleware {
|
||||
return &AuthMiddleware{
|
||||
userRepo: userRepo,
|
||||
logger: logs.NewLogger(),
|
||||
}
|
||||
}
|
||||
|
||||
// getJWTSecret 获取JWT密钥
|
||||
func (m *AuthMiddleware) getJWTSecret() []byte {
|
||||
// 在实际项目中,应该从配置文件或环境变量中读取
|
||||
secret := os.Getenv("JWT_SECRET")
|
||||
if secret == "" {
|
||||
secret = "pig-farm-controller-secret-key" // 默认密钥
|
||||
}
|
||||
return []byte(secret)
|
||||
}
|
||||
|
||||
// GenerateToken 为用户生成JWT token
|
||||
func (m *AuthMiddleware) GenerateToken(userID uint, username string) (string, error) {
|
||||
claims := JWTClaims{
|
||||
UserID: userID,
|
||||
Username: username,
|
||||
RegisteredClaims: jwt.RegisteredClaims{
|
||||
ExpiresAt: jwt.NewNumericDate(time.Now().Add(24 * time.Hour)), // 24小时过期
|
||||
IssuedAt: jwt.NewNumericDate(time.Now()),
|
||||
NotBefore: jwt.NewNumericDate(time.Now()),
|
||||
Issuer: "pig-farm-controller",
|
||||
},
|
||||
}
|
||||
|
||||
token := jwt.NewWithClaims(jwt.SigningMethodHS256, claims)
|
||||
return token.SignedString(m.getJWTSecret())
|
||||
}
|
||||
|
||||
// Handle 鉴权中间件处理函数
|
||||
func (m *AuthMiddleware) Handle() gin.HandlerFunc {
|
||||
return func(c *gin.Context) {
|
||||
// 从请求头中获取认证信息
|
||||
authHeader := c.GetHeader("Authorization")
|
||||
if authHeader == "" {
|
||||
c.JSON(http.StatusUnauthorized, gin.H{"error": "缺少认证信息"})
|
||||
c.Abort()
|
||||
return
|
||||
}
|
||||
|
||||
// 检查Bearer token格式
|
||||
if !strings.HasPrefix(authHeader, "Bearer ") {
|
||||
c.JSON(http.StatusUnauthorized, gin.H{"error": "认证信息格式错误"})
|
||||
c.Abort()
|
||||
return
|
||||
}
|
||||
|
||||
// 解析token
|
||||
tokenString := strings.TrimPrefix(authHeader, "Bearer ")
|
||||
|
||||
// 验证token并获取用户信息
|
||||
user, err := m.getUserFromJWT(tokenString)
|
||||
if err != nil {
|
||||
if err == gorm.ErrRecordNotFound {
|
||||
c.JSON(http.StatusUnauthorized, gin.H{"error": "用户不存在"})
|
||||
} else {
|
||||
m.logger.Error("Token验证失败: " + err.Error())
|
||||
c.JSON(http.StatusUnauthorized, gin.H{"error": "无效的认证令牌"})
|
||||
}
|
||||
c.Abort()
|
||||
return
|
||||
}
|
||||
|
||||
// 将用户信息保存到上下文中,供后续处理函数使用
|
||||
c.Set("user", user)
|
||||
|
||||
// 继续处理请求
|
||||
c.Next()
|
||||
}
|
||||
}
|
||||
|
||||
// getUserFromJWT 从JWT token中获取用户信息
|
||||
func (m *AuthMiddleware) getUserFromJWT(tokenString string) (*AuthUser, error) {
|
||||
// 解析token
|
||||
token, err := jwt.ParseWithClaims(tokenString, &JWTClaims{}, func(token *jwt.Token) (interface{}, error) {
|
||||
return m.getJWTSecret(), nil
|
||||
})
|
||||
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
// 验证token
|
||||
if !token.Valid {
|
||||
return nil, gorm.ErrRecordNotFound
|
||||
}
|
||||
|
||||
// 获取声明
|
||||
claims, ok := token.Claims.(*JWTClaims)
|
||||
if !ok {
|
||||
return nil, gorm.ErrRecordNotFound
|
||||
}
|
||||
|
||||
// 根据用户ID查找用户
|
||||
userModel, err := m.userRepo.FindByID(claims.UserID)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
user := &AuthUser{
|
||||
ID: userModel.ID,
|
||||
Username: userModel.Username,
|
||||
}
|
||||
|
||||
return user, nil
|
||||
}
|
||||
@@ -77,12 +77,12 @@ func (c *Config) Load(path string) error {
|
||||
// 读取配置文件
|
||||
data, err := os.ReadFile(path)
|
||||
if err != nil {
|
||||
return fmt.Errorf("failed to read config file: %v", err)
|
||||
return fmt.Errorf("配置文件读取失败: %v", err)
|
||||
}
|
||||
|
||||
// 解析YAML配置
|
||||
if err := yaml.Unmarshal(data, c); err != nil {
|
||||
return fmt.Errorf("failed to parse config file: %v", err)
|
||||
return fmt.Errorf("配置文件解析失败: %v", err)
|
||||
}
|
||||
|
||||
return nil
|
||||
|
||||
126
internal/controller/user/user.go
Normal file
126
internal/controller/user/user.go
Normal file
@@ -0,0 +1,126 @@
|
||||
// Package user 提供用户相关功能的控制器
|
||||
// 实现用户注册、登录等操作
|
||||
package user
|
||||
|
||||
import (
|
||||
"net/http"
|
||||
|
||||
"git.huangwc.com/pig/pig-farm-controller/internal/api/middleware"
|
||||
"git.huangwc.com/pig/pig-farm-controller/internal/logs"
|
||||
"git.huangwc.com/pig/pig-farm-controller/internal/storage/repository"
|
||||
"github.com/gin-gonic/gin"
|
||||
"golang.org/x/crypto/bcrypt"
|
||||
)
|
||||
|
||||
// Controller 用户控制器
|
||||
type Controller struct {
|
||||
userRepo repository.UserRepo
|
||||
logger *logs.Logger
|
||||
}
|
||||
|
||||
// NewController 创建用户控制器实例
|
||||
func NewController(userRepo repository.UserRepo) *Controller {
|
||||
return &Controller{
|
||||
userRepo: userRepo,
|
||||
logger: logs.NewLogger(),
|
||||
}
|
||||
}
|
||||
|
||||
// RegisterRequest 注册请求结构体
|
||||
type RegisterRequest struct {
|
||||
Username string `json:"username" binding:"required"`
|
||||
Password string `json:"password" binding:"required"`
|
||||
}
|
||||
|
||||
// RegisterResponse 注册响应结构体
|
||||
type RegisterResponse struct {
|
||||
ID uint `json:"id"`
|
||||
Username string `json:"username"`
|
||||
CreatedAt string `json:"created_at"`
|
||||
}
|
||||
|
||||
// Register 用户注册
|
||||
func (c *Controller) Register(ctx *gin.Context) {
|
||||
var req RegisterRequest
|
||||
if err := ctx.ShouldBindJSON(&req); err != nil {
|
||||
ctx.JSON(http.StatusBadRequest, gin.H{"error": "请求参数错误"})
|
||||
return
|
||||
}
|
||||
|
||||
user, err := c.userRepo.CreateUser(req.Username, req.Password)
|
||||
if err != nil {
|
||||
c.logger.Error("创建用户失败: " + err.Error())
|
||||
ctx.JSON(http.StatusInternalServerError, gin.H{"error": "创建用户失败"})
|
||||
return
|
||||
}
|
||||
|
||||
response := RegisterResponse{
|
||||
ID: user.ID,
|
||||
Username: user.Username,
|
||||
CreatedAt: user.CreatedAt.Format("2006-01-02 15:04:05"),
|
||||
}
|
||||
|
||||
ctx.JSON(http.StatusOK, gin.H{
|
||||
"message": "用户创建成功",
|
||||
"user": response,
|
||||
})
|
||||
}
|
||||
|
||||
// LoginRequest 登录请求结构体
|
||||
type LoginRequest struct {
|
||||
Username string `json:"username" binding:"required"`
|
||||
Password string `json:"password" binding:"required"`
|
||||
}
|
||||
|
||||
// LoginResponse 登录响应结构体
|
||||
type LoginResponse struct {
|
||||
ID uint `json:"id"`
|
||||
Username string `json:"username"`
|
||||
Token string `json:"token"`
|
||||
CreatedAt string `json:"created_at"`
|
||||
}
|
||||
|
||||
// Login 用户登录
|
||||
func (c *Controller) Login(ctx *gin.Context) {
|
||||
var req LoginRequest
|
||||
if err := ctx.ShouldBindJSON(&req); err != nil {
|
||||
ctx.JSON(http.StatusBadRequest, gin.H{"error": "请求参数错误"})
|
||||
return
|
||||
}
|
||||
|
||||
// 查找用户
|
||||
user, err := c.userRepo.FindByUsername(req.Username)
|
||||
if err != nil {
|
||||
c.logger.Error("查找用户失败: " + err.Error())
|
||||
ctx.JSON(http.StatusUnauthorized, gin.H{"error": "用户名或密码错误"})
|
||||
return
|
||||
}
|
||||
|
||||
// 验证密码
|
||||
err = bcrypt.CompareHashAndPassword([]byte(user.PasswordHash), []byte(req.Password))
|
||||
if err != nil {
|
||||
ctx.JSON(http.StatusUnauthorized, gin.H{"error": "用户名或密码错误"})
|
||||
return
|
||||
}
|
||||
|
||||
// 生成JWT访问令牌
|
||||
authMiddleware := middleware.NewAuthMiddleware(c.userRepo)
|
||||
token, err := authMiddleware.GenerateToken(user.ID, user.Username)
|
||||
if err != nil {
|
||||
c.logger.Error("生成令牌失败: " + err.Error())
|
||||
ctx.JSON(http.StatusInternalServerError, gin.H{"error": "登录失败"})
|
||||
return
|
||||
}
|
||||
|
||||
response := LoginResponse{
|
||||
ID: user.ID,
|
||||
Username: user.Username,
|
||||
Token: token,
|
||||
CreatedAt: user.CreatedAt.Format("2006-01-02 15:04:05"),
|
||||
}
|
||||
|
||||
ctx.JSON(http.StatusOK, gin.H{
|
||||
"message": "登录成功",
|
||||
"user": response,
|
||||
})
|
||||
}
|
||||
@@ -10,6 +10,7 @@ import (
|
||||
"git.huangwc.com/pig/pig-farm-controller/internal/config"
|
||||
"git.huangwc.com/pig/pig-farm-controller/internal/logs"
|
||||
"git.huangwc.com/pig/pig-farm-controller/internal/storage/db"
|
||||
"git.huangwc.com/pig/pig-farm-controller/internal/storage/repository"
|
||||
"git.huangwc.com/pig/pig-farm-controller/internal/task"
|
||||
)
|
||||
|
||||
@@ -25,6 +26,9 @@ type Application struct {
|
||||
// TaskExecutor 任务执行器组件实例
|
||||
TaskExecutor *task.Executor
|
||||
|
||||
// UserRepo 用户仓库实例
|
||||
UserRepo repository.UserRepo
|
||||
|
||||
// Config 应用配置
|
||||
Config *config.Config
|
||||
|
||||
@@ -46,8 +50,11 @@ func NewApplication(cfg *config.Config) *Application {
|
||||
// 初始化存储组件
|
||||
store := db.NewStorage(connectionString, maxOpenConns, maxIdleConns, connMaxLifetime)
|
||||
|
||||
// 初始化用户仓库
|
||||
userRepo := repository.NewUserRepo(store.GetDB())
|
||||
|
||||
// 初始化API组件
|
||||
apiInstance := api.NewAPI(cfg)
|
||||
apiInstance := api.NewAPI(cfg, userRepo)
|
||||
|
||||
// 初始化任务执行器组件(使用5个工作协程)
|
||||
taskExecutor := task.NewExecutor(5)
|
||||
@@ -56,6 +63,7 @@ func NewApplication(cfg *config.Config) *Application {
|
||||
Storage: store,
|
||||
API: apiInstance,
|
||||
TaskExecutor: taskExecutor,
|
||||
UserRepo: userRepo,
|
||||
Config: cfg,
|
||||
logger: logs.NewLogger(),
|
||||
}
|
||||
@@ -66,19 +74,19 @@ func NewApplication(cfg *config.Config) *Application {
|
||||
func (app *Application) Start() error {
|
||||
// 启动存储组件
|
||||
if err := app.Storage.Connect(); err != nil {
|
||||
return fmt.Errorf("failed to connect to storage: %v", err)
|
||||
return fmt.Errorf("存储连接失败: %v", err)
|
||||
}
|
||||
app.logger.Info("Storage connected successfully")
|
||||
app.logger.Info("存储连接成功")
|
||||
|
||||
// 启动API组件
|
||||
if err := app.API.Start(); err != nil {
|
||||
return fmt.Errorf("failed to start API: %v", err)
|
||||
return fmt.Errorf("API启动失败: %v", err)
|
||||
}
|
||||
app.logger.Info("API started successfully")
|
||||
app.logger.Info("API启动成功")
|
||||
|
||||
// 启动任务执行器组件
|
||||
app.TaskExecutor.Start()
|
||||
app.logger.Info("Task executor started successfully")
|
||||
app.logger.Info("任务执行器启动成功")
|
||||
|
||||
return nil
|
||||
}
|
||||
@@ -88,18 +96,18 @@ func (app *Application) Start() error {
|
||||
func (app *Application) Stop() error {
|
||||
// 停止API组件
|
||||
if err := app.API.Stop(); err != nil {
|
||||
app.logger.Error(fmt.Sprintf("Failed to stop API: %v", err))
|
||||
app.logger.Error(fmt.Sprintf("API停止失败: %v", err))
|
||||
}
|
||||
|
||||
// 停止任务执行器组件
|
||||
app.TaskExecutor.Stop()
|
||||
app.logger.Info("Task executor stopped successfully")
|
||||
app.logger.Info("任务执行器已停止")
|
||||
|
||||
// 停止存储组件
|
||||
if err := app.Storage.Disconnect(); err != nil {
|
||||
return fmt.Errorf("failed to disconnect from storage: %v", err)
|
||||
return fmt.Errorf("存储断开连接失败: %v", err)
|
||||
}
|
||||
app.logger.Info("Storage disconnected successfully")
|
||||
app.logger.Info("存储断开连接成功")
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -23,20 +23,20 @@ func NewLogger() *Logger {
|
||||
|
||||
// Info 记录信息级别日志
|
||||
func (l *Logger) Info(message string) {
|
||||
l.logger.Printf("[INFO] %s %s", time.Now().Format(time.RFC3339), message)
|
||||
l.logger.Printf("[信息] %s %s", time.Now().Format(time.RFC3339), message)
|
||||
}
|
||||
|
||||
// Error 记录错误级别日志
|
||||
func (l *Logger) Error(message string) {
|
||||
l.logger.Printf("[ERROR] %s %s", time.Now().Format(time.RFC3339), message)
|
||||
l.logger.Printf("[错误] %s %s", time.Now().Format(time.RFC3339), message)
|
||||
}
|
||||
|
||||
// Debug 记录调试级别日志
|
||||
func (l *Logger) Debug(message string) {
|
||||
l.logger.Printf("[DEBUG] %s %s", time.Now().Format(time.RFC3339), message)
|
||||
l.logger.Printf("[调试] %s %s", time.Now().Format(time.RFC3339), message)
|
||||
}
|
||||
|
||||
// Warn 记录警告级别日志
|
||||
func (l *Logger) Warn(message string) {
|
||||
l.logger.Printf("[WARN] %s %s", time.Now().Format(time.RFC3339), message)
|
||||
l.logger.Printf("[警告] %s %s", time.Now().Format(time.RFC3339), message)
|
||||
}
|
||||
|
||||
35
internal/model/user.go
Normal file
35
internal/model/user.go
Normal file
@@ -0,0 +1,35 @@
|
||||
// Package model 提供数据模型定义
|
||||
// 包含用户、猪舍、饲料等相关数据结构
|
||||
package model
|
||||
|
||||
import (
|
||||
"time"
|
||||
|
||||
"gorm.io/gorm"
|
||||
)
|
||||
|
||||
// User 代表系统用户
|
||||
type User struct {
|
||||
// ID 用户ID
|
||||
ID uint `gorm:"primaryKey;column:id" json:"id"`
|
||||
|
||||
// Username 用户名
|
||||
Username string `gorm:"uniqueIndex;not null;column:username" json:"username"`
|
||||
|
||||
// PasswordHash 密码哈希值
|
||||
PasswordHash string `gorm:"not null;column:password_hash" json:"-"`
|
||||
|
||||
// CreatedAt 创建时间
|
||||
CreatedAt time.Time `gorm:"column:created_at" json:"created_at"`
|
||||
|
||||
// UpdatedAt 更新时间
|
||||
UpdatedAt time.Time `gorm:"column:updated_at" json:"updated_at"`
|
||||
|
||||
// DeletedAt 删除时间(用于软删除)
|
||||
DeletedAt gorm.DeletedAt `gorm:"index;column:deleted_at" json:"-"`
|
||||
}
|
||||
|
||||
// TableName 指定User模型对应的数据库表名
|
||||
func (User) TableName() string {
|
||||
return "users"
|
||||
}
|
||||
@@ -49,25 +49,25 @@ func NewPostgresStorage(connectionString string, maxOpenConns, maxIdleConns, con
|
||||
// Connect 建立与PostgreSQL数据库的连接
|
||||
// 使用GORM建立数据库连接
|
||||
func (ps *PostgresStorage) Connect() error {
|
||||
ps.logger.Info("Connecting to PostgreSQL database")
|
||||
ps.logger.Info("正在连接PostgreSQL数据库")
|
||||
|
||||
var err error
|
||||
ps.db, err = gorm.Open(postgres.Open(ps.connectionString), &gorm.Config{})
|
||||
if err != nil {
|
||||
ps.logger.Error(fmt.Sprintf("Failed to connect to database: %v", err))
|
||||
return fmt.Errorf("failed to connect to database: %v", err)
|
||||
ps.logger.Error(fmt.Sprintf("数据库连接失败: %v", err))
|
||||
return fmt.Errorf("数据库连接失败: %v", err)
|
||||
}
|
||||
|
||||
// 测试连接
|
||||
sqlDB, err := ps.db.DB()
|
||||
if err != nil {
|
||||
ps.logger.Error(fmt.Sprintf("Failed to get database instance: %v", err))
|
||||
return fmt.Errorf("failed to get database instance: %v", err)
|
||||
ps.logger.Error(fmt.Sprintf("获取数据库实例失败: %v", err))
|
||||
return fmt.Errorf("获取数据库实例失败: %v", err)
|
||||
}
|
||||
|
||||
if err = sqlDB.Ping(); err != nil {
|
||||
ps.logger.Error(fmt.Sprintf("Failed to ping database: %v", err))
|
||||
return fmt.Errorf("failed to ping database: %v", err)
|
||||
ps.logger.Error(fmt.Sprintf("数据库连接测试失败: %v", err))
|
||||
return fmt.Errorf("数据库连接测试失败: %v", err)
|
||||
}
|
||||
|
||||
// 设置连接池参数
|
||||
@@ -75,7 +75,7 @@ func (ps *PostgresStorage) Connect() error {
|
||||
sqlDB.SetMaxIdleConns(ps.maxIdleConns)
|
||||
sqlDB.SetConnMaxLifetime(time.Duration(ps.connMaxLifetime) * time.Second)
|
||||
|
||||
ps.logger.Info("Successfully connected to PostgreSQL database")
|
||||
ps.logger.Info("PostgreSQL数据库连接成功")
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -83,19 +83,19 @@ func (ps *PostgresStorage) Connect() error {
|
||||
// 安全地关闭所有数据库连接
|
||||
func (ps *PostgresStorage) Disconnect() error {
|
||||
if ps.db != nil {
|
||||
ps.logger.Info("Disconnecting from PostgreSQL database")
|
||||
ps.logger.Info("正在断开PostgreSQL数据库连接")
|
||||
|
||||
sqlDB, err := ps.db.DB()
|
||||
if err != nil {
|
||||
ps.logger.Error(fmt.Sprintf("Failed to get database instance: %v", err))
|
||||
return fmt.Errorf("failed to get database instance: %v", err)
|
||||
ps.logger.Error(fmt.Sprintf("获取数据库实例失败: %v", err))
|
||||
return fmt.Errorf("获取数据库实例失败: %v", err)
|
||||
}
|
||||
|
||||
if err := sqlDB.Close(); err != nil {
|
||||
ps.logger.Error(fmt.Sprintf("Failed to close database connection: %v", err))
|
||||
return fmt.Errorf("failed to close database connection: %v", err)
|
||||
ps.logger.Error(fmt.Sprintf("关闭数据库连接失败: %v", err))
|
||||
return fmt.Errorf("关闭数据库连接失败: %v", err)
|
||||
}
|
||||
ps.logger.Info("Successfully disconnected from PostgreSQL database")
|
||||
ps.logger.Info("PostgreSQL数据库连接已断开")
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
84
internal/storage/repository/user.go
Normal file
84
internal/storage/repository/user.go
Normal file
@@ -0,0 +1,84 @@
|
||||
// Package repository 提供数据访问层实现
|
||||
// 包含各种数据实体的仓库接口和实现
|
||||
package repository
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
|
||||
"git.huangwc.com/pig/pig-farm-controller/internal/model"
|
||||
"golang.org/x/crypto/bcrypt"
|
||||
"gorm.io/gorm"
|
||||
)
|
||||
|
||||
// UserRepo 用户仓库接口
|
||||
type UserRepo interface {
|
||||
// CreateUser 创建新用户
|
||||
CreateUser(username, password string) (*model.User, error)
|
||||
|
||||
// FindByUsername 根据用户名查找用户
|
||||
FindByUsername(username string) (*model.User, error)
|
||||
|
||||
// FindByID 根据ID查找用户
|
||||
FindByID(id uint) (*model.User, error)
|
||||
}
|
||||
|
||||
// userRepo 用户仓库实现
|
||||
type userRepo struct {
|
||||
db *gorm.DB
|
||||
}
|
||||
|
||||
// NewUserRepo 创建用户仓库实例
|
||||
func NewUserRepo(db *gorm.DB) UserRepo {
|
||||
return &userRepo{
|
||||
db: db,
|
||||
}
|
||||
}
|
||||
|
||||
// CreateUser 创建新用户
|
||||
func (r *userRepo) CreateUser(username, password string) (*model.User, error) {
|
||||
// 检查用户是否已存在
|
||||
var existingUser model.User
|
||||
result := r.db.Where("username = ?", username).First(&existingUser)
|
||||
if result.Error == nil {
|
||||
return nil, fmt.Errorf("用户已存在")
|
||||
}
|
||||
|
||||
// 对密码进行哈希处理
|
||||
hashedPassword, err := bcrypt.GenerateFromPassword([]byte(password), bcrypt.DefaultCost)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("密码加密失败: %v", err)
|
||||
}
|
||||
|
||||
// 创建新用户
|
||||
user := &model.User{
|
||||
Username: username,
|
||||
PasswordHash: string(hashedPassword),
|
||||
}
|
||||
|
||||
result = r.db.Create(user)
|
||||
if result.Error != nil {
|
||||
return nil, fmt.Errorf("用户创建失败: %v", result.Error)
|
||||
}
|
||||
|
||||
return user, nil
|
||||
}
|
||||
|
||||
// FindByUsername 根据用户名查找用户
|
||||
func (r *userRepo) FindByUsername(username string) (*model.User, error) {
|
||||
var user model.User
|
||||
result := r.db.Where("username = ?", username).First(&user)
|
||||
if result.Error != nil {
|
||||
return nil, result.Error
|
||||
}
|
||||
return &user, nil
|
||||
}
|
||||
|
||||
// FindByID 根据ID查找用户
|
||||
func (r *userRepo) FindByID(id uint) (*model.User, error) {
|
||||
var user model.User
|
||||
result := r.db.First(&user, id)
|
||||
if result.Error != nil {
|
||||
return nil, result.Error
|
||||
}
|
||||
return &user, nil
|
||||
}
|
||||
@@ -65,7 +65,7 @@ func (tq *TaskQueue) AddTask(task Task) {
|
||||
priority: task.GetPriority(),
|
||||
}
|
||||
heap.Push(tq.queue, item)
|
||||
tq.logger.Info("Task added to queue: " + task.GetID())
|
||||
tq.logger.Info("任务已添加到队列: " + task.GetID())
|
||||
}
|
||||
|
||||
// GetNextTask 获取下一个要执行的任务(优先级最高的任务)
|
||||
@@ -79,7 +79,7 @@ func (tq *TaskQueue) GetNextTask() Task {
|
||||
|
||||
// 获取优先级最高的任务
|
||||
item := heap.Pop(tq.queue).(*taskItem)
|
||||
tq.logger.Info("Task retrieved from queue: " + item.task.GetID())
|
||||
tq.logger.Info("从队列中获取任务: " + item.task.GetID())
|
||||
return item.task
|
||||
}
|
||||
|
||||
@@ -160,7 +160,7 @@ func NewExecutor(workers int) *Executor {
|
||||
|
||||
// Start 启动任务执行器
|
||||
func (e *Executor) Start() {
|
||||
e.logger.Info(fmt.Sprintf("Starting task executor with %d workers", e.workers))
|
||||
e.logger.Info(fmt.Sprintf("正在启动任务执行器,工作协程数: %d", e.workers))
|
||||
|
||||
// 启动工作协程
|
||||
for i := 0; i < e.workers; i++ {
|
||||
@@ -168,12 +168,12 @@ func (e *Executor) Start() {
|
||||
go e.worker(i)
|
||||
}
|
||||
|
||||
e.logger.Info("Task executor started successfully")
|
||||
e.logger.Info("任务执行器启动成功")
|
||||
}
|
||||
|
||||
// Stop 停止任务执行器
|
||||
func (e *Executor) Stop() {
|
||||
e.logger.Info("Stopping task executor")
|
||||
e.logger.Info("正在停止任务执行器")
|
||||
|
||||
// 取消上下文
|
||||
e.cancel()
|
||||
@@ -181,37 +181,37 @@ func (e *Executor) Stop() {
|
||||
// 等待所有工作协程结束
|
||||
e.wg.Wait()
|
||||
|
||||
e.logger.Info("Task executor stopped successfully")
|
||||
e.logger.Info("任务执行器已停止")
|
||||
}
|
||||
|
||||
// SubmitTask 提交任务到执行器
|
||||
func (e *Executor) SubmitTask(task Task) {
|
||||
e.taskQueue.AddTask(task)
|
||||
e.logger.Info("Task submitted: " + task.GetID())
|
||||
e.logger.Info("任务已提交: " + task.GetID())
|
||||
}
|
||||
|
||||
// worker 工作协程
|
||||
func (e *Executor) worker(id int) {
|
||||
defer e.wg.Done()
|
||||
|
||||
e.logger.Info(fmt.Sprintf("Worker (id = %d) started", id))
|
||||
e.logger.Info(fmt.Sprintf("工作协程(id = %d)已启动", id))
|
||||
|
||||
for {
|
||||
select {
|
||||
case <-e.ctx.Done():
|
||||
e.logger.Info(fmt.Sprintf("Worker %d stopped", id))
|
||||
e.logger.Info(fmt.Sprintf("工作协程 %d 已停止", id))
|
||||
return
|
||||
default:
|
||||
// 获取下一个任务
|
||||
task := e.taskQueue.GetNextTask()
|
||||
if task != nil {
|
||||
e.logger.Info(fmt.Sprintf("Worker %d executing task: %s", id, task.GetID()))
|
||||
e.logger.Info(fmt.Sprintf("工作协程 %d 正在执行任务: %s", id, task.GetID()))
|
||||
|
||||
// 执行任务
|
||||
if err := task.Execute(); err != nil {
|
||||
e.logger.Error("Task execution failed: " + task.GetID() + ", error: " + err.Error())
|
||||
e.logger.Error("任务执行失败: " + task.GetID() + ", 错误: " + err.Error())
|
||||
} else {
|
||||
e.logger.Info("Task executed successfully: " + task.GetID())
|
||||
e.logger.Info("任务执行成功: " + task.GetID())
|
||||
}
|
||||
} else {
|
||||
// 没有任务时短暂休眠
|
||||
|
||||
Reference in New Issue
Block a user