package service import ( "errors" "fmt" "strings" "hy2xs-admin/dao" "hy2xs-admin/model/bo" "hy2xs-admin/model/constant" "hy2xs-admin/model/dto" "hy2xs-admin/model/entity" "hy2xs-admin/model/vo" "hy2xs-admin/util" ) func PagePeer(peerPageDto dto.PeerPageDto) ([]vo.PeerVo, int64, error) { peers, total, err := dao.PagePeer(peerPageDto) if err != nil { return nil, 0, err } onlineUsers, _ := Hysteria2Online() result := make([]vo.PeerVo, 0, len(peers)) for _, p := range peers { item := vo.PeerVo{ BaseVo: vo.BaseVo{Id: *p.Id, CreateTime: *p.CreateTime}, Name: strVal(p.Name), Remark: strVal(p.Remark), AuthId: strVal(p.AuthId), QuotaBytes: int64Val(p.QuotaBytes), DownloadBytes: int64Val(p.DownloadBytes), UploadBytes: int64Val(p.UploadBytes), ExpiresAt: int64Val(p.ExpiresAt), MaxDevices: int64Val(p.MaxDevices), Disabled: int64Val(p.Disabled), BannedUntil: int64Val(p.BannedUntil), LastConnectionAt: int64Val(p.LastConnectionAt), } authID := strVal(p.AuthId) if v, ok := onlineUsers[authID]; ok { item.Online = true item.OnlineDevices = v } result = append(result, item) } return result, total, nil } func CreatePeer(peerDto dto.PeerSaveDto) (vo.PeerVo, error) { if peerDto.Name == nil || *peerDto.Name == "" { return vo.PeerVo{}, errors.New(constant.InvalidError) } if ExistPeerName(*peerDto.Name, 0) { return vo.PeerVo{}, errors.New(fmt.Sprintf("name %s already exists", *peerDto.Name)) } secret := "" if peerDto.Secret != nil && *peerDto.Secret != "" { secret = *peerDto.Secret } else { generated, err := util.RandomString(24) if err != nil { return vo.PeerVo{}, err } secret = fmt.Sprintf("%s.%s", *peerDto.Name, generated) } authId, err := util.RandomString(18) if err != nil { return vo.PeerVo{}, err } secretDigest, err := PeerSecretDigest(secret) if err != nil { return vo.PeerVo{}, err } secretEncrypted, err := EncryptPeerSecret(secret) if err != nil { return vo.PeerVo{}, err } peer := entity.Peer{ Name: peerDto.Name, Remark: peerDto.Remark, AuthId: &authId, SecretDigest: &secretDigest, SecretEncrypted: &secretEncrypted, QuotaBytes: peerDto.QuotaBytes, ExpiresAt: peerDto.ExpiresAt, MaxDevices: peerDto.MaxDevices, Disabled: peerDto.Disabled, } id, saveErr := dao.SavePeer(peer) if saveErr != nil { return vo.PeerVo{}, saveErr } return GetPeerVo(id) } func UpdatePeer(id int64, peerDto dto.PeerUpdateDto) error { updates := map[string]interface{}{} if peerDto.Name != nil && *peerDto.Name != "" { updates["name"] = *peerDto.Name } if peerDto.Secret != nil && *peerDto.Secret != "" { digest, err := PeerSecretDigest(*peerDto.Secret) if err != nil { return err } enc, err := EncryptPeerSecret(*peerDto.Secret) if err != nil { return err } updates["secret_digest"] = digest updates["secret_ciphertext"] = enc } if peerDto.QuotaBytes != nil { updates["quota_bytes"] = *peerDto.QuotaBytes } if peerDto.ExpiresAt != nil { updates["expires_at"] = *peerDto.ExpiresAt } if peerDto.MaxDevices != nil { updates["max_devices"] = *peerDto.MaxDevices } if peerDto.Disabled != nil { updates["disabled"] = *peerDto.Disabled } if peerDto.Remark != nil { updates["remark"] = *peerDto.Remark } return dao.UpdatePeer([]int64{id}, updates) } func DeletePeer(id int64) error { return dao.DeletePeer([]int64{id}) } func GetPeerVo(id int64) (vo.PeerVo, error) { p, err := dao.GetPeer("id = ?", id) if err != nil { return vo.PeerVo{}, err } return vo.PeerVo{ BaseVo: vo.BaseVo{Id: *p.Id, CreateTime: *p.CreateTime}, Name: strVal(p.Name), Remark: strVal(p.Remark), AuthId: strVal(p.AuthId), QuotaBytes: int64Val(p.QuotaBytes), DownloadBytes: int64Val(p.DownloadBytes), UploadBytes: int64Val(p.UploadBytes), ExpiresAt: int64Val(p.ExpiresAt), MaxDevices: int64Val(p.MaxDevices), Disabled: int64Val(p.Disabled), BannedUntil: int64Val(p.BannedUntil), LastConnectionAt: int64Val(p.LastConnectionAt), }, nil } func ResetPeerTraffic(id int64) error { return dao.UpdatePeer([]int64{id}, map[string]interface{}{"download_bytes": 0, "upload_bytes": 0}) } func ReleaseKickPeer(id int64) error { return dao.UpdatePeer([]int64{id}, map[string]interface{}{"banned_until": 0}) } func KickPeer(id int64, bannedUntil int64) error { if err := dao.UpdatePeer([]int64{id}, map[string]interface{}{"banned_until": bannedUntil}); err != nil { return err } return Hysteria2Kick([]int64{id}, bannedUntil) } func BuildPeerClientConfig(id int64) (vo.PeerClientConfigVo, error) { url, err := Hysteria2Url(id) if err != nil { return vo.PeerClientConfigVo{}, err } return vo.PeerClientConfigVo{Url: url}, nil } func ListExportPeer(includeSecrets bool) ([]bo.PeerExport, error) { peers, err := dao.ListPeer("1=1") if err != nil { return nil, errors.New(constant.SysError) } out := make([]bo.PeerExport, 0, len(peers)) for _, item := range peers { ex := bo.PeerExport{ Id: int64Val(item.Id), AuthId: strVal(item.AuthId), Name: strVal(item.Name), Remark: strVal(item.Remark), QuotaBytes: int64Val(item.QuotaBytes), DownloadBytes: int64Val(item.DownloadBytes), UploadBytes: int64Val(item.UploadBytes), ExpiresAt: int64Val(item.ExpiresAt), MaxDevices: int64Val(item.MaxDevices), Disabled: int64Val(item.Disabled), BannedUntil: int64Val(item.BannedUntil), LastConnectionAt: int64Val(item.LastConnectionAt), } if includeSecrets && item.SecretEncrypted != nil { if dec, derr := DecryptPeerSecret(*item.SecretEncrypted); derr == nil { ex.Secret = dec } } out = append(out, ex) } return out, nil } // preparedPeerImport — запись импорта со всем криптоматериалом, посчитанным // заранее. // // Крипто выносится ИЗ транзакции сознательно. PeerSecretDigest и // EncryptPeerSecret читают ключи из таблицы `config`, то есть ходят в ту же // базу; делать это, удерживая открытую запись, значит без нужды держать // блокировку на время AES по каждой из тысяч записей. Внутри транзакции должна // остаться только работа с таблицей пиров. type preparedPeerImport struct { source bo.PeerExport name string authID string remark string quota int64 expires int64 maxDevices int64 disabled int64 // Задан, только если секрет пришёл в файле: у существующего пира секрет // перезаписывается лишь в этом случае. hasExplicitSecret bool explicitDigest string explicitCipher string // Готовятся всегда: понадобятся, если запись окажется новой. createDigest string createCipher string createAuthID string } func preparePeerImport(items []bo.PeerExport) ([]preparedPeerImport, error) { prepared := make([]preparedPeerImport, 0, len(items)) for _, item := range items { name := strings.TrimSpace(item.Name) authID := strings.TrimSpace(item.AuthId) maxDevices := item.MaxDevices if maxDevices <= 0 { maxDevices = 3 } entry := preparedPeerImport{ source: item, name: name, authID: authID, remark: item.Remark, quota: item.QuotaBytes, expires: item.ExpiresAt, maxDevices: maxDevices, disabled: item.Disabled, } explicitSecret := strings.TrimSpace(item.Secret) if explicitSecret != "" { digest, err := PeerSecretDigest(explicitSecret) if err != nil { return nil, err } cipher, err := EncryptPeerSecret(explicitSecret) if err != nil { return nil, err } entry.hasExplicitSecret = true entry.explicitDigest = digest entry.explicitCipher = cipher entry.createDigest = digest entry.createCipher = cipher } else { generated, err := util.RandomString(24) if err != nil { return nil, err } createSecret := fmt.Sprintf("%s.%s", name, generated) digest, err := PeerSecretDigest(createSecret) if err != nil { return nil, err } cipher, err := EncryptPeerSecret(createSecret) if err != nil { return nil, err } entry.createDigest = digest entry.createCipher = cipher } entry.createAuthID = authID if entry.createAuthID == "" { generated, err := util.RandomString(18) if err != nil { return nil, err } entry.createAuthID = generated } prepared = append(prepared, entry) } return prepared, nil } // UpsertPeerExport применяет выгрузку пиров целиком или не применяет вовсе. // // Три прохода, и каждый отвечает за своё: // // 1. ValidatePeerImportBatch — содержимое файла, без обращения к базе; // 2. preparePeerImport — весь криптоматериал, без обращения к таблице пиров; // 3. одна транзакция — только записи. // // Раньше третьего прохода не существовало: записи шли по одной, каждая своим // оператором. Комментарий обещал «либо целиком, либо никак», но UNIQUE-конфликт // на 37-й записи оставлял 36 применённых, и откатить это оператор уже не мог. // Конфликт не гипотетический: пусть в базе есть A(auth_id=a, name=alice) и // B(auth_id=b, name=bob), а файл несёт (auth_id=a, name=bob). Поиск найдёт A // по auth_id и переименует его в bob — прямо в UNIQUE(name). func UpsertPeerExport(items []bo.PeerExport) error { if err := ValidatePeerImportBatch(items); err != nil { return err } prepared, err := preparePeerImport(items) if err != nil { return err } return dao.WithPeerTx(func(tx dao.PeerTx) error { for _, entry := range prepared { if err := applyPeerImportEntry(tx, entry); err != nil { return err } } return nil }) } func applyPeerImportEntry(tx dao.PeerTx, entry preparedPeerImport) error { var existing entity.Peer var err error if entry.authID != "" { existing, err = tx.GetPeer("auth_id = ?", entry.authID) } if err != nil || existing.Id == nil { existing, err = tx.GetPeer("name = ?", entry.name) } // Пир установщика не переопределяется импортом ни при каком совпадении: // его секрет живёт ещё и в /etc/hy2xs/bootstrap-admin.secret. if err == nil && existing.Name != nil && *existing.Name == ReservedBootstrapPeerName { return fmt.Errorf( "peer import: пир %q принадлежит установщику и не может быть изменён импортом", ReservedBootstrapPeerName, ) } if err == nil && existing.Id != nil { updates := map[string]interface{}{ "name": entry.name, "remark": entry.remark, "quota_bytes": entry.quota, "download_bytes": entry.source.DownloadBytes, "upload_bytes": entry.source.UploadBytes, "expires_at": entry.expires, "max_devices": entry.maxDevices, "disabled": entry.disabled, "banned_until": entry.source.BannedUntil, "last_connection_at": entry.source.LastConnectionAt, } if entry.authID != "" { updates["auth_id"] = entry.authID } if entry.hasExplicitSecret { updates["secret_digest"] = entry.explicitDigest updates["secret_ciphertext"] = entry.explicitCipher } return tx.UpdatePeer([]int64{*existing.Id}, updates) } name := entry.name remark := entry.remark authID := entry.createAuthID digest := entry.createDigest cipher := entry.createCipher quota := entry.quota expires := entry.expires maxDevices := entry.maxDevices disabled := entry.disabled download := entry.source.DownloadBytes upload := entry.source.UploadBytes bannedUntil := entry.source.BannedUntil lastConnection := entry.source.LastConnectionAt peer := entity.Peer{ Name: &name, Remark: &remark, AuthId: &authID, SecretDigest: &digest, SecretEncrypted: &cipher, QuotaBytes: "a, DownloadBytes: &download, UploadBytes: &upload, ExpiresAt: &expires, MaxDevices: &maxDevices, Disabled: &disabled, BannedUntil: &bannedUntil, LastConnectionAt: &lastConnection, } _, saveErr := tx.SavePeer(peer) return saveErr } func ExistPeerName(name string, id int64) bool { var err error if id != 0 { _, err = dao.GetPeer("name = ? and id != ?", name, id) } else { _, err = dao.GetPeer("name = ?", name) } return err == nil } func UpdatePeerLastConnectionAt(id int64, conAt int64) error { return dao.UpdatePeer([]int64{id}, map[string]interface{}{"last_connection_at": conAt}) } func strVal(v *string) string { if v == nil { return "" } return *v } func int64Val(v *int64) int64 { if v == nil { return 0 } return *v }