166 lines
5.1 KiB
Go
166 lines
5.1 KiB
Go
package main
|
||
|
||
import (
|
||
"context"
|
||
"crypto/sha1"
|
||
"database/sql"
|
||
"encoding/hex"
|
||
"encoding/json"
|
||
"flag"
|
||
"fmt"
|
||
"log"
|
||
"strings"
|
||
"time"
|
||
|
||
_ "github.com/go-sql-driver/mysql"
|
||
"google.golang.org/grpc"
|
||
"google.golang.org/grpc/credentials/insecure"
|
||
roomv1 "hyapp.local/api/proto/room/v1"
|
||
"hyapp/pkg/appcode"
|
||
)
|
||
|
||
type roomCandidate struct {
|
||
RoomID string `json:"room_id"`
|
||
RoomShortID string `json:"room_short_id"`
|
||
OwnerUserID int64 `json:"owner_user_id"`
|
||
VisibleRegionID int64 `json:"visible_region_id"`
|
||
OwnerRegionID int64 `json:"owner_region_id"`
|
||
}
|
||
|
||
func main() {
|
||
var roomDSN string
|
||
var userDSN string
|
||
var roomGRPCAddr string
|
||
var app string
|
||
var shortID string
|
||
var limit int
|
||
var apply bool
|
||
var adminID uint64
|
||
flag.StringVar(&roomDSN, "room_mysql_dsn", "", "room-service MySQL DSN")
|
||
flag.StringVar(&userDSN, "user_mysql_dsn", "", "user-service MySQL DSN")
|
||
flag.StringVar(&roomGRPCAddr, "room_grpc_addr", "", "room-service gRPC address; required with -apply")
|
||
flag.StringVar(&app, "app_code", appcode.Default, "app_code scope")
|
||
flag.StringVar(&shortID, "room_short_id", "", "optional room_short_id filter")
|
||
flag.IntVar(&limit, "limit", 500, "maximum rooms to scan")
|
||
flag.BoolVar(&apply, "apply", false, "call room-service AdminUpdateRoom; omitted means dry-run")
|
||
flag.Uint64Var(&adminID, "admin_id", 0, "admin id written to AdminUpdateRoom")
|
||
flag.Parse()
|
||
|
||
if strings.TrimSpace(roomDSN) == "" || strings.TrimSpace(userDSN) == "" {
|
||
log.Fatal("room_mysql_dsn and user_mysql_dsn are required")
|
||
}
|
||
if apply && strings.TrimSpace(roomGRPCAddr) == "" {
|
||
log.Fatal("room_grpc_addr is required when -apply is set")
|
||
}
|
||
ctx, cancel := context.WithTimeout(context.Background(), 2*time.Minute)
|
||
defer cancel()
|
||
|
||
roomDB, err := sql.Open("mysql", roomDSN)
|
||
if err != nil {
|
||
log.Fatalf("open room mysql failed: %v", err)
|
||
}
|
||
defer roomDB.Close()
|
||
userDB, err := sql.Open("mysql", userDSN)
|
||
if err != nil {
|
||
log.Fatalf("open user mysql failed: %v", err)
|
||
}
|
||
defer userDB.Close()
|
||
|
||
mismatches, err := findMismatches(ctx, roomDB, userDB, appcode.Normalize(app), shortID, limit)
|
||
if err != nil {
|
||
log.Fatalf("find mismatches failed: %v", err)
|
||
}
|
||
encoder := json.NewEncoder(log.Writer())
|
||
for _, item := range mismatches {
|
||
if err := encoder.Encode(item); err != nil {
|
||
log.Fatalf("write dry-run row failed: %v", err)
|
||
}
|
||
}
|
||
log.Printf("mismatch_count=%d apply=%t", len(mismatches), apply)
|
||
if !apply || len(mismatches) == 0 {
|
||
return
|
||
}
|
||
|
||
conn, err := grpc.DialContext(ctx, roomGRPCAddr, grpc.WithTransportCredentials(insecure.NewCredentials()))
|
||
if err != nil {
|
||
log.Fatalf("dial room-service failed: %v", err)
|
||
}
|
||
defer conn.Close()
|
||
client := roomv1.NewRoomCommandServiceClient(conn)
|
||
for _, item := range mismatches {
|
||
// 迁移必须走 Room Cell 命令链路,让 command log、snapshot、rooms 和 room_list_entries 在同一事务内收敛。
|
||
if _, err := client.AdminUpdateRoom(ctx, &roomv1.AdminUpdateRoomRequest{
|
||
Meta: &roomv1.RequestMeta{
|
||
AppCode: appcode.Normalize(app),
|
||
// 后台迁移仍须携带真实房主作为命令操作者;Room Cell 依赖该字段完成审计和完整性校验。
|
||
ActorUserId: item.OwnerUserID,
|
||
RequestId: "room-visible-region-migration",
|
||
CommandId: migrationCommandID(item),
|
||
RoomId: item.RoomID,
|
||
SentAtMs: time.Now().UTC().UnixMilli(),
|
||
},
|
||
VisibleRegionId: &item.OwnerRegionID,
|
||
AdminId: adminID,
|
||
AdminName: "room-visible-region-migration",
|
||
}); err != nil {
|
||
log.Fatalf("migrate room_id=%s short_id=%s failed: %v", item.RoomID, item.RoomShortID, err)
|
||
}
|
||
}
|
||
log.Printf("migrated_count=%d", len(mismatches))
|
||
}
|
||
|
||
func findMismatches(ctx context.Context, roomDB *sql.DB, userDB *sql.DB, app string, shortID string, limit int) ([]roomCandidate, error) {
|
||
if limit <= 0 {
|
||
limit = 500
|
||
}
|
||
where := "WHERE app_code = ? AND status <> 'deleted'"
|
||
args := []any{app}
|
||
if strings.TrimSpace(shortID) != "" {
|
||
where += " AND room_short_id = ?"
|
||
args = append(args, strings.TrimSpace(shortID))
|
||
}
|
||
args = append(args, limit)
|
||
rows, err := roomDB.QueryContext(ctx, `
|
||
SELECT room_id, room_short_id, owner_user_id, COALESCE(visible_region_id, 0)
|
||
FROM rooms
|
||
`+where+`
|
||
ORDER BY room_id ASC
|
||
LIMIT ?
|
||
`, args...)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
defer rows.Close()
|
||
|
||
mismatches := make([]roomCandidate, 0)
|
||
for rows.Next() {
|
||
var item roomCandidate
|
||
if err := rows.Scan(&item.RoomID, &item.RoomShortID, &item.OwnerUserID, &item.VisibleRegionID); err != nil {
|
||
return nil, err
|
||
}
|
||
var ownerRegionID int64
|
||
err := userDB.QueryRowContext(ctx, `
|
||
SELECT COALESCE(region_id, 0)
|
||
FROM users
|
||
WHERE app_code = ? AND user_id = ? AND COALESCE(region_id, 0) > 0
|
||
`, app, item.OwnerUserID).Scan(&ownerRegionID)
|
||
if err == sql.ErrNoRows {
|
||
continue
|
||
}
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
if ownerRegionID == item.VisibleRegionID {
|
||
continue
|
||
}
|
||
item.OwnerRegionID = ownerRegionID
|
||
mismatches = append(mismatches, item)
|
||
}
|
||
return mismatches, rows.Err()
|
||
}
|
||
|
||
func migrationCommandID(item roomCandidate) string {
|
||
sum := sha1.Sum([]byte(fmt.Sprintf("%s:%d:%d", item.RoomID, item.VisibleRegionID, item.OwnerRegionID)))
|
||
return "room-region-migrate:" + hex.EncodeToString(sum[:])
|
||
}
|