refactor(core): completely decommission legacy internal/apps, internal/platform/bootstrap, and internal/router/v1

This commit is contained in:
ryan
2026-08-28 07:28:45 +08:00
parent df3c5ae756
commit dc49b72c29
182 changed files with 946 additions and 19604 deletions
+16 -265
View File
@@ -5,33 +5,26 @@ package cmd
import (
"context"
"errors"
"fmt"
"log"
"net"
"net/http"
"sync"
"time"
"github.com/Rain-kl/Wavelet/core"
"github.com/Rain-kl/Wavelet/core/extpoints"
gwrunner "github.com/Rain-kl/Wavelet/internal/apps/message_gateway/runner"
"github.com/Rain-kl/Wavelet/internal/infra/config"
"github.com/Rain-kl/Wavelet/internal/infra/task/scheduler"
"github.com/Rain-kl/Wavelet/internal/infra/task/worker"
"github.com/Rain-kl/Wavelet/internal/platform/bootstrap"
"github.com/Rain-kl/Wavelet/internal/router"
"github.com/Rain-kl/Wavelet/pkg/util"
"github.com/Rain-kl/Wavelet/plugins/domain/admin"
"github.com/Rain-kl/Wavelet/plugins/domain/auth"
"github.com/Rain-kl/Wavelet/plugins/domain/cap"
"github.com/Rain-kl/Wavelet/plugins/domain/message_gateway"
"github.com/Rain-kl/Wavelet/plugins/domain/risk_control"
"github.com/Rain-kl/Wavelet/plugins/domain/system"
"github.com/Rain-kl/Wavelet/plugins/domain/upload"
"github.com/Rain-kl/Wavelet/plugins/domain/user"
"github.com/Rain-kl/Wavelet/plugins/drivers/driver_asynq_cron"
"github.com/Rain-kl/Wavelet/plugins/drivers/driver_asynq_worker"
"github.com/Rain-kl/Wavelet/plugins/drivers/driver_http"
"github.com/Rain-kl/Wavelet/plugins/infra/cache"
"github.com/Rain-kl/Wavelet/plugins/infra/database"
"github.com/Rain-kl/Wavelet/plugins/infra/logger"
"github.com/Rain-kl/Wavelet/plugins/infra/storage"
"github.com/hibiken/asynq"
)
// newWaveletApp creates a core.App wired with Wavelet platform infrastructure, domain plugins, and profile drivers.
@@ -43,7 +36,7 @@ func newWaveletApp(profile core.Profile) *core.App {
core.WithShutdownTimeout(time.Duration(config.Config.App.GracefulShutdownTimeout)*time.Second),
)
// Register standard infrastructure plugins
// 1. Register standard infrastructure plugins
app.Use(
database.New(),
cache.New(),
@@ -51,272 +44,30 @@ func newWaveletApp(profile core.Profile) *core.App {
storage.New(),
)
// Register domain plugins
// 2. Register all 8 domain business plugins
app.Use(
auth.New(),
user.New(),
message_gateway.New(),
risk_control.New(),
admin.New(),
upload.New(),
cap.New(),
system.New(),
)
// Bind Goose migration runner
// 3. Bind Goose migration runner
app.SetMigrationRunner(func(_ context.Context, _ []extpoints.MigrationEntry) error {
runMigrations()
return nil
})
// Mount drivers for each aspect
// 4. Mount runtime drivers for each aspect
app.Use(
newWaveletHTTPDriver(profile),
newWaveletWorkerDriver(profile),
newWaveletSchedulerDriver(profile),
driver_http.New(driver_http.WithAddr(config.Config.App.Addr)),
driver_asynq_worker.New(),
driver_asynq_cron.New(),
)
return app
}
type waveletHTTPDriver struct {
profile core.Profile
server *http.Server
mu sync.Mutex
running bool
}
func newWaveletHTTPDriver(profile core.Profile) *waveletHTTPDriver {
return &waveletHTTPDriver{profile: profile}
}
func (d *waveletHTTPDriver) Name() string {
return "driver_wavelet_http"
}
func (d *waveletHTTPDriver) Apply(ctx *core.Context) error {
return ctx.RegisterDriver(d)
}
func (d *waveletHTTPDriver) Type() core.DriverType {
return core.DriverTypeHTTP
}
//nolint:contextcheck
func (d *waveletHTTPDriver) Start(ctx context.Context) error {
d.mu.Lock()
defer d.mu.Unlock()
if d.running {
return nil
}
bootstrap.RegisterAPI()
runBootstrap(bootstrap.Options{API: true})
engine, err := router.BuildEngine()
if err != nil {
return fmt.Errorf("[API] build router engine failed: %w", err)
}
srv := &http.Server{
Addr: config.Config.App.Addr,
Handler: engine,
ReadHeaderTimeout: 10 * time.Second,
}
listener, err := (&net.ListenConfig{}).Listen(ctx, "tcp", config.Config.App.Addr)
if err != nil {
return fmt.Errorf("[API] listen on %s failed: %w", config.Config.App.Addr, err)
}
mode := "API"
if d.profile == core.ProfileAll {
mode = "API + Worker + Scheduler"
}
printStartupBanner(startupState{
mode: mode,
relationalDB: latestMigrationState.relationalDB,
clickHouseDB: latestMigrationState.clickHouseDB,
listensForHTTP: true,
})
d.server = srv
d.running = true
util.Go(func() {
log.Printf("[API] server listening on %s\n", config.Config.App.Addr)
if serveErr := srv.Serve(listener); serveErr != nil && !errors.Is(serveErr, http.ErrServerClosed) {
log.Fatalf("[API] server failed: %v\n", serveErr)
}
})
return nil
}
func (d *waveletHTTPDriver) Stop(ctx context.Context) error {
d.mu.Lock()
defer d.mu.Unlock()
if !d.running {
return nil
}
d.running = false
var err error
if d.server != nil {
err = d.server.Shutdown(ctx)
d.server = nil
}
bootstrap.Stop(ctx)
log.Println("[API] server exited")
return err
}
type waveletWorkerDriver struct {
profile core.Profile
server *asynq.Server
mu sync.Mutex
running bool
}
func newWaveletWorkerDriver(profile core.Profile) *waveletWorkerDriver {
return &waveletWorkerDriver{profile: profile}
}
func (d *waveletWorkerDriver) Name() string {
return "driver_wavelet_worker"
}
func (d *waveletWorkerDriver) Apply(ctx *core.Context) error {
return ctx.RegisterDriver(d)
}
func (d *waveletWorkerDriver) Type() core.DriverType {
return core.DriverTypeWorker
}
//nolint:contextcheck
func (d *waveletWorkerDriver) Start(_ context.Context) error {
d.mu.Lock()
defer d.mu.Unlock()
if d.running {
return nil
}
if d.profile == core.ProfileAll {
bootstrap.RegisterAll()
} else {
bootstrap.RegisterWorker()
}
runBootstrap(bootstrap.Options{})
if d.profile == core.ProfileWorker {
printStartupBanner(startupState{
mode: "Worker",
relationalDB: latestMigrationState.relationalDB,
clickHouseDB: latestMigrationState.clickHouseDB,
})
}
util.Go(func() {
if err := gwrunner.Start(context.Background()); err != nil {
log.Printf("[Worker] message gateway stopped: %v", err)
}
})
log.Println("[Worker] 启动任务处理服务")
srv, err := worker.StartWorkerServer()
if err != nil {
return fmt.Errorf("[Worker] 启动失败: %w", err)
}
d.server = srv
d.running = true
return nil
}
func (d *waveletWorkerDriver) Stop(_ context.Context) error {
d.mu.Lock()
defer d.mu.Unlock()
if !d.running {
return nil
}
d.running = false
if d.server != nil {
d.server.Stop()
d.server.Shutdown()
d.server = nil
}
log.Println("[Worker] 任务处理服务已退出")
return nil
}
type waveletSchedulerDriver struct {
profile core.Profile
mu sync.Mutex
running bool
}
func newWaveletSchedulerDriver(profile core.Profile) *waveletSchedulerDriver {
return &waveletSchedulerDriver{profile: profile}
}
func (d *waveletSchedulerDriver) Name() string {
return "driver_wavelet_scheduler"
}
func (d *waveletSchedulerDriver) Apply(ctx *core.Context) error {
return ctx.RegisterDriver(d)
}
func (d *waveletSchedulerDriver) Type() core.DriverType {
return core.DriverTypeScheduler
}
//nolint:contextcheck
func (d *waveletSchedulerDriver) Start(_ context.Context) error {
d.mu.Lock()
defer d.mu.Unlock()
if d.running {
return nil
}
if d.profile == core.ProfileAll {
bootstrap.RegisterAll()
} else {
bootstrap.RegisterScheduler()
}
runBootstrap(bootstrap.Options{})
if d.profile == core.ProfileSchedule {
printStartupBanner(startupState{
mode: "Scheduler",
relationalDB: latestMigrationState.relationalDB,
clickHouseDB: latestMigrationState.clickHouseDB,
})
}
log.Println("[Scheduler] 启动定时任务调度服务")
if err := scheduler.ReloadScheduler(); err != nil {
return fmt.Errorf("[Scheduler] 启动失败: %w", err)
}
d.running = true
return nil
}
func (d *waveletSchedulerDriver) Stop(_ context.Context) error {
d.mu.Lock()
defer d.mu.Unlock()
if !d.running {
return nil
}
d.running = false
scheduler.StopScheduler()
log.Println("[Scheduler] 定时任务调度服务已退出")
return nil
}
+14 -5
View File
@@ -26,9 +26,9 @@ func TestNewWaveletAppProfiles(t *testing.T) {
require.NotNil(t, app)
assert.Equal(t, prof, app.Profile())
// Verify 4 infra plugins + 5 domain plugins + 3 driver plugins registered
// Verify 4 infra plugins + 8 domain plugins + 3 driver plugins registered
plugins := app.Plugins()
assert.Len(t, plugins, 12)
assert.Len(t, plugins, 15)
// Verify each standard infra plugin is registered
_, ok := app.Plugin("database")
@@ -59,14 +59,23 @@ func TestNewWaveletAppProfiles(t *testing.T) {
_, ok = app.Plugin("admin")
assert.True(t, ok, "admin plugin missing")
_, ok = app.Plugin("upload")
assert.True(t, ok, "upload plugin missing")
_, ok = app.Plugin("cap")
assert.True(t, ok, "cap plugin missing")
_, ok = app.Plugin("system")
assert.True(t, ok, "system plugin missing")
// Verify driver plugins
_, ok = app.Plugin("driver_wavelet_http")
_, ok = app.Plugin("driver_http")
assert.True(t, ok, "http driver missing")
_, ok = app.Plugin("driver_wavelet_worker")
_, ok = app.Plugin("driver_asynq_worker")
assert.True(t, ok, "worker driver missing")
_, ok = app.Plugin("driver_wavelet_scheduler")
_, ok = app.Plugin("driver_asynq_cron")
assert.True(t, ok, "scheduler driver missing")
})
}
+4 -3
View File
@@ -1,6 +1,6 @@
// Copyright 2026 Arctel.net
// SPDX-License-Identifier: Apache-2.0
// Package cmd provides CLI command entry points.
//
//nolint:unused
package cmd
import (
@@ -14,6 +14,7 @@ import (
"github.com/Rain-kl/Wavelet/internal/infra/persistence/migrator"
)
//nolint:unused // startup banner formatting utilities
type startupState struct {
mode string
relationalDB migrator.Report
-17
View File
@@ -1,17 +0,0 @@
// Copyright 2026 Arctel.net
// SPDX-License-Identifier: Apache-2.0
package cmd
import (
"context"
"github.com/Rain-kl/Wavelet/internal/platform/bootstrap"
"github.com/Rain-kl/Wavelet/pkg/trace"
)
func runBootstrap(opts bootstrap.Options) {
ctx, span := trace.Start(context.Background(), "bootstrap.Init")
defer span.End()
bootstrap.Init(ctx, opts)
}
+3 -5
View File
@@ -13,12 +13,11 @@ import (
"os"
"strings"
"github.com/Rain-kl/Wavelet/internal/apps/oauth"
"github.com/Rain-kl/Wavelet/internal/infra/persistence"
"github.com/Rain-kl/Wavelet/internal/infra/persistence/migrator"
"github.com/Rain-kl/Wavelet/internal/model"
"github.com/Rain-kl/Wavelet/internal/platform/bootstrap"
"github.com/Rain-kl/Wavelet/internal/repository"
"github.com/Rain-kl/Wavelet/plugins/domain/auth"
"github.com/spf13/cobra"
"gorm.io/gorm"
)
@@ -52,7 +51,6 @@ var resetPasswdCmd = &cobra.Command{
},
Run: func(_ *cobra.Command, _ []string) {
ctx := context.Background()
runBootstrap(bootstrap.Options{})
var username string
if usernameFlag != "" {
@@ -101,7 +99,7 @@ var resetPasswdCmd = &cobra.Command{
var tokens []model.AccessToken
if err := tx.Where("user_id = ?", user.ID).Find(&tokens).Error; err == nil {
for _, token := range tokens {
oauth.InvalidateCachedToken(ctx, token.TokenHash)
auth.InvalidateCachedToken(ctx, token.TokenHash)
}
}
@@ -111,7 +109,7 @@ var resetPasswdCmd = &cobra.Command{
log.Fatalf("重置密码失败: %v\n", err)
}
oauth.InvalidateCachedUser(ctx, user.ID)
auth.InvalidateCachedUser(ctx, user.ID)
fmt.Println("成功重置密码!")
fmt.Printf("用户名: %s\n", user.Username)