package main import ( "context" "fmt" "sync/atomic" "time" "go_win_template/internal/bus" "go_win_template/internal/cache" "go_win_template/internal/plugins" "go_win_template/internal/service" "go_win_template/internal/worker" ) // App 主应用结构体 type App struct { ctx context.Context users []User nextID atomic.Int64 userCache *cache.Cache eventBus *bus.EventBus batchWorker *service.BatchWorker taskPool *worker.Pool requestCount atomic.Int64 paymentSvc paymentService } // paymentService 支付服务接口(定义在 app 层,不与 payment 包直接耦合) // payment 插件在启动时通过类型断言注入此接口 type paymentService interface { PayLogin() (map[string]any, error) PayCheckAuth() (map[string]any, error) PayRefresh() (map[string]any, error) PayGetAccountInfo() (map[string]any, error) PayQueryOrder(orderID string) (map[string]any, error) PayQueryOrderByOutNo(outTradeNo string) (map[string]any, error) PayRequestRefund(req map[string]any) (map[string]any, error) PayGetBill(page, pageSize int) (map[string]any, error) PayGetBillByDate(startDate, endDate string) (map[string]any, error) } // Response 统一响应 type Response struct { Code int `json:"code"` Message string `json:"message"` Data interface{} `json:"data,omitempty"` } // User 用户模型 type User struct { ID int64 `json:"id"` Name string `json:"name"` Email string `json:"email"` CreatedAt string `json:"created_at"` } // ───────────────────────────────────────────── // 初始化 // ───────────────────────────────────────────── func NewApp() *App { return &App{ users: make([]User, 0), userCache: cache.NewCache(cache.WithTTL(10*time.Minute), cache.WithCleanupInterval(1*time.Minute)), eventBus: bus.NewEventBus(), batchWorker: service.NewBatchWorker(8), taskPool: worker.NewPool(worker.Options{Workers: 4, BufSize: 500}), } } func (a *App) Startup(ctx context.Context) { a.ctx = ctx fmt.Println("[app] startup complete") // 注入支付服务 a.injectPaymentService() } func (a *App) injectPaymentService() { payP := plugins.Get("payment") if payP == nil { fmt.Println("[app] payment plugin not loaded (use -tags payment to enable)") return } if getter, ok := payP.(interface{ GetPaymentService() any }); ok { if svc, ok := getter.GetPaymentService().(paymentService); ok { a.paymentSvc = svc fmt.Println("[app] payment service injected successfully") return } } fmt.Println("[app] payment plugin found but service type mismatch") } // ───────────────────────────────────────────── // IPC 接口 — 基础功能 // ───────────────────────────────────────────── func (a *App) GetUserList() []User { a.requestCount.Add(1) if items, ok := a.userCache.Get("user_list"); ok { if cached, ok := items.([]User); ok { return cached } } return a.users } func (a *App) CreateUser(name, email string) *Response { if name == "" || email == "" { return &Response{Code: 400, Message: "名称和邮箱不能为空"} } id := a.nextID.Add(1) now := time.Now().Format("2006-01-02 15:04:05") user := User{ID: id, Name: name, Email: email, CreatedAt: now} a.users = append(a.users, user) a.userCache.Set("user_list", a.users) a.eventBus.PublishAsync("user:created", user) return &Response{Code: 0, Message: "success", Data: user} } func (a *App) ClearUsers() int { count := len(a.users) a.users = make([]User, 0) a.userCache.Delete("user_list") return count } func (a *App) Stats() map[string]interface{} { running, total, failed := a.taskPool.Stats() br, bt, bs, bf := a.batchWorker.Stats() return map[string]interface{}{ "total_users": len(a.users), "cache_items": a.userCache.Len(), "requests": a.requestCount.Load(), "pool_running": running, "pool_total": total, "pool_failed": failed, "batch_running": br, "batch_total": bt, "batch_success": bs, "batch_failed": bf, "plugins": plugins.List(), "plugin_count": plugins.Count(), "payment_loaded": a.paymentSvc != nil, } } // ───────────────────────────────────────────── // IPC 接口 — 支付功能 // ───────────────────────────────────────────── func (a *App) PayLogin() *Response { if a.paymentSvc == nil { return &Response{Code: 503, Message: "支付插件未加载,使用 -tags payment 构建或检查 plugin.json"} } data, err := a.paymentSvc.PayLogin() if err != nil { return &Response{Code: 500, Message: err.Error()} } return &Response{Code: 0, Message: asString(data, "message"), Data: data} } func (a *App) PayCheckAuth() *Response { if a.paymentSvc == nil { return &Response{Code: 503, Message: "支付插件未加载"} } data, err := a.paymentSvc.PayCheckAuth() if err != nil { return &Response{Code: 500, Message: err.Error()} } return &Response{Code: 0, Message: asString(data, "message"), Data: data} } func (a *App) PayRefreshToken() *Response { if a.paymentSvc == nil { return &Response{Code: 503, Message: "支付插件未加载"} } data, err := a.paymentSvc.PayRefresh() if err != nil { return &Response{Code: 500, Message: err.Error()} } return &Response{Code: 0, Message: asString(data, "message"), Data: data} } func (a *App) PayGetAccountInfo() *Response { if a.paymentSvc == nil { return &Response{Code: 503, Message: "支付插件未加载"} } data, err := a.paymentSvc.PayGetAccountInfo() if err != nil { return &Response{Code: 500, Message: err.Error()} } return &Response{Code: 0, Message: asString(data, "message"), Data: data} } func (a *App) PayQueryOrder(orderID string) *Response { if a.paymentSvc == nil { return &Response{Code: 503, Message: "支付插件未加载"} } data, err := a.paymentSvc.PayQueryOrder(orderID) if err != nil { return &Response{Code: 500, Message: err.Error()} } return &Response{Code: 0, Message: asString(data, "message"), Data: data} } func (a *App) PayQueryOrderByOutNo(outTradeNo string) *Response { if a.paymentSvc == nil { return &Response{Code: 503, Message: "支付插件未加载"} } data, err := a.paymentSvc.PayQueryOrderByOutNo(outTradeNo) if err != nil { return &Response{Code: 500, Message: err.Error()} } return &Response{Code: 0, Message: asString(data, "message"), Data: data} } func (a *App) PayRequestRefund(orderID string, refundAmt float64, reason string) *Response { if a.paymentSvc == nil { return &Response{Code: 503, Message: "支付插件未加载"} } data, err := a.paymentSvc.PayRequestRefund(map[string]any{ "order_id": orderID, "refund_amt": refundAmt, "reason": reason, }) if err != nil { return &Response{Code: 500, Message: err.Error()} } return &Response{Code: 0, Message: asString(data, "message"), Data: data} } func (a *App) PayGetBill(page, pageSize int) *Response { if a.paymentSvc == nil { return &Response{Code: 503, Message: "支付插件未加载"} } data, err := a.paymentSvc.PayGetBill(page, pageSize) if err != nil { return &Response{Code: 500, Message: err.Error()} } return &Response{Code: 0, Message: asString(data, "message"), Data: data} } func (a *App) PayGetBillByDate(startDate, endDate string) *Response { if a.paymentSvc == nil { return &Response{Code: 503, Message: "支付插件未加载"} } data, err := a.paymentSvc.PayGetBillByDate(startDate, endDate) if err != nil { return &Response{Code: 500, Message: err.Error()} } return &Response{Code: 0, Message: asString(data, "message"), Data: data} } // ───────────────────────────────────────────── // 高并发批量处理 // ───────────────────────────────────────────── func (a *App) BatchCreateUsers(names, emails []string) *Response { if len(names) != len(emails) || len(names) == 0 { return &Response{Code: 400, Message: "names 和 emails 长度需一致且不为空"} } results := a.batchWorker.ProcessConcurrent(a.ctx, createBatchItems(names, emails), func(ctx context.Context, item any) (any, error) { bp := item.(batchItem) id := a.nextID.Add(1) user := User{ID: id, Name: bp.name, Email: bp.email, CreatedAt: time.Now().Format("2006-01-02 15:04:05")} a.users = append(a.users, user) return user, nil }) a.userCache.Set("user_list", a.users) a.eventBus.PublishAsync("users:batch_created", results) return &Response{ Code: 0, Message: fmt.Sprintf("成功 %d 条,失败 %d 条", len(results.Success), len(results.Fail)), Data: map[string]any{"success": len(results.Success), "failed": len(results.Fail)}, } } func (a *App) SubmitTask(taskName, taskData string) *Response { job := &exampleJob{name: taskName, data: taskData} if !a.taskPool.Submit(job) { return &Response{Code: 503, Message: "任务队列已满,请稍后重试"} } return &Response{Code: 0, Message: "任务已入队", Data: map[string]any{"task": taskName}} } // ───────────────────────────────────────────── // 内部类型 // ───────────────────────────────────────────── type batchItem struct{ name, email string } func createBatchItems(names, emails []string) []any { items := make([]any, len(names)) for i := range names { items[i] = batchItem{name: names[i], email: emails[i]} } return items } type exampleJob struct{ name, data string } func (j *exampleJob) JobID() string { return j.name } func (j *exampleJob) Execute(ctx context.Context) error { select { case <-ctx.Done(): return ctx.Err() case <-time.After(100 * time.Millisecond): return nil } } // asString 安全地从 map 中提取字符串字段 func asString(m map[string]any, key string) string { if v, ok := m[key].(string); ok { return v } return "" }