package environment import ( "context" "crypto/rand" "database/sql" "encoding/hex" "encoding/json" "errors" "net" "regexp" "strings" "time" ) var exitIDPattern = regexp.MustCompile(`^[A-Za-z0-9][A-Za-z0-9._/-]{0,127}$`) var nodeIDPattern = regexp.MustCompile(`^[A-Za-z0-9][A-Za-z0-9._:-]{0,127}$`) type NetworkExit struct { ID string `json:"id"` Protocol string `json:"protocol"` Host string `json:"host"` Port int `json:"port"` Username string `json:"username"` Password string `json:"password"` ExpectedPublicIP string `json:"expected_public_ip,omitempty"` ExpectedRegion string `json:"expected_region,omitempty"` ObservedPublicIP string `json:"observed_public_ip,omitempty"` ObservedRegion string `json:"observed_region,omitempty"` HealthStatus string `json:"health_status"` LastCheckReason string `json:"last_check_reason,omitempty"` Version int64 `json:"version"` LastCheckedAt *time.Time `json:"last_checked_at,omitempty"` CreatedAt time.Time `json:"created_at"` UpdatedAt time.Time `json:"updated_at"` } // NetworkExitAccess is the internal runtime view of a persisted network exit. type NetworkExitAccess struct { NetworkExit } type ExitObservation struct { PublicIP string Region string } type EnvironmentContext struct { Env AccountID string `json:"account_id"` AccountStatus string `json:"account_status"` BindingVersion int64 `json:"binding_version"` RuntimeCleanupPending bool `json:"runtime_cleanup_pending,omitempty"` RuntimeCleanupBindingVersion int64 `json:"runtime_cleanup_binding_version,omitempty"` RuntimeCleanupRuntimeID string `json:"runtime_cleanup_runtime_id,omitempty"` RuntimeCleanupNetworkID string `json:"runtime_cleanup_network_id,omitempty"` Exit NetworkExit `json:"network_exit"` RuntimeID string `json:"runtime_id,omitempty"` RuntimeNetworkID string `json:"runtime_network_id,omitempty"` RuntimeNodeID string `json:"runtime_node_id,omitempty"` } type EnvironmentAction struct { OperationID string Action string AccountID string BrowserEnvAlias string NetworkExitID string BindingVersion int64 Outcome string ReasonCode string } func (s *Store) CreateNetworkExit(ctx context.Context, exit NetworkExit) (NetworkExit, error) { exit.ID = "exit-" + newHubID() exit.Protocol, exit.Host = strings.ToLower(strings.TrimSpace(exit.Protocol)), strings.TrimSpace(exit.Host) exit.ExpectedPublicIP, exit.ExpectedRegion = strings.TrimSpace(exit.ExpectedPublicIP), strings.TrimSpace(exit.ExpectedRegion) if !validNetworkExit(exit) { return NetworkExit{}, ErrInvalid } row := s.db.QueryRowContext(ctx, ` INSERT INTO network_exit (exit_id, protocol, host, port, username, password, expected_public_ip, expected_region) VALUES ($1, $2, $3, $4, $5, $6, NULLIF($7, '')::inet, $8) RETURNING exit_id`, exit.ID, exit.Protocol, exit.Host, exit.Port, exit.Username, exit.Password, exit.ExpectedPublicIP, exit.ExpectedRegion) if err := row.Scan(&exit.ID); err != nil { return NetworkExit{}, publicDatabaseError(err) } return s.GetNetworkExit(ctx, exit.ID) } func validNetworkExit(exit NetworkExit) bool { if exit.Protocol != "http" && exit.Protocol != "https" && exit.Protocol != "socks4" && exit.Protocol != "socks5" { return false } if !validExitHost(exit.Host) || exit.Port < 1 || exit.Port > 65535 || !validExitCredentials(exit.Username, exit.Password) { return false } if exit.ExpectedPublicIP != "" && net.ParseIP(exit.ExpectedPublicIP) == nil { return false } return validOptionalRegion(exit.ExpectedRegion) } func validExitCredentials(username, password string) bool { if len(username) > 255 || len(password) > 255 || (username == "" && password != "") { return false } for _, value := range username + password { if value < 0x20 || value == 0x7f { return false } } return true } func validExitHost(host string) bool { if host == "" || len(host) > 253 || strings.ContainsAny(host, "@/[]?# \t\r\n") { return false } if net.ParseIP(host) != nil { return true } if strings.HasPrefix(host, ".") || strings.HasSuffix(host, ".") || strings.Contains(host, "..") { return false } for _, label := range strings.Split(host, ".") { if len(label) > 63 || strings.HasPrefix(label, "-") || strings.HasSuffix(label, "-") { return false } for _, character := range label { if (character < 'a' || character > 'z') && (character < 'A' || character > 'Z') && (character < '0' || character > '9') && character != '-' { return false } } } return true } func validOptionalRegion(region string) bool { if len(region) > 64 { return false } for _, character := range region { if character < 0x20 || character == 0x7f { return false } } return true } func (s *Store) UpdateNetworkExit(ctx context.Context, id string, input NetworkExit) (NetworkExit, error) { if !exitIDPattern.MatchString(id) { return NetworkExit{}, ErrInvalid } input.Protocol, input.Host = strings.ToLower(strings.TrimSpace(input.Protocol)), strings.TrimSpace(input.Host) input.ExpectedPublicIP, input.ExpectedRegion = strings.TrimSpace(input.ExpectedPublicIP), strings.TrimSpace(input.ExpectedRegion) if !validNetworkExit(input) { return NetworkExit{}, ErrInvalid } tx, err := s.db.BeginTx(ctx, nil) if err != nil { return NetworkExit{}, errors.New("begin network exit update") } defer tx.Rollback() var active bool if err := tx.QueryRowContext(ctx, `SELECT EXISTS (SELECT 1 FROM browser_env b JOIN network_exit n ON n.id = b.exit_id WHERE n.exit_id=$1 AND b.runtime_id IS NOT NULL AND b.runtime_lease_until > now())`, id).Scan(&active); err != nil { return NetworkExit{}, errors.New("check network exit activity") } if active { return NetworkExit{}, ErrConflict } if _, err := tx.ExecContext(ctx, `UPDATE network_exit SET protocol=$2,host=$3,port=$4,username=$5,password=$6,expected_public_ip=NULLIF($7,'')::inet,expected_region=$8,health_status=CASE WHEN health_status='disabled' THEN 'disabled' ELSE 'unchecked' END,last_check_reason='exit_configuration_changed',version=version+1,updated_at=now() WHERE exit_id=$1`, id, input.Protocol, input.Host, input.Port, input.Username, input.Password, input.ExpectedPublicIP, input.ExpectedRegion); err != nil { return NetworkExit{}, rowError(err) } if err := tx.Commit(); err != nil { return NetworkExit{}, err } return s.GetNetworkExit(ctx, id) } func (s *Store) EnableNetworkExit(ctx context.Context, id string) (NetworkExit, error) { if !exitIDPattern.MatchString(id) { return NetworkExit{}, ErrInvalid } if _, err := s.db.ExecContext(ctx, `UPDATE network_exit SET health_status='unchecked',last_check_reason='exit_enabled',version=version+1,updated_at=now() WHERE exit_id=$1 AND health_status='disabled'`, id); err != nil { return NetworkExit{}, rowError(err) } return s.GetNetworkExit(ctx, id) } func (s *Store) DeleteNetworkExit(ctx context.Context, id string) error { if !exitIDPattern.MatchString(id) { return ErrInvalid } tx, err := s.db.BeginTx(ctx, nil) if err != nil { return errors.New("begin network exit delete") } defer tx.Rollback() var used bool if err := tx.QueryRowContext(ctx, `SELECT EXISTS (SELECT 1 FROM browser_env b JOIN network_exit n ON n.id = b.exit_id WHERE n.exit_id=$1)`, id).Scan(&used); err != nil { return errors.New("check network exit bindings") } if used { return ErrConflict } result, err := tx.ExecContext(ctx, `DELETE FROM network_exit WHERE exit_id=$1`, id) if err != nil { return rowError(err) } count, err := result.RowsAffected() if err != nil { return err } if count != 1 { return ErrNotFound } return tx.Commit() } func (s *Store) ListNetworkExits(ctx context.Context) ([]NetworkExit, error) { rows, err := s.db.QueryContext(ctx, networkExitSelect+` ORDER BY network.created_at, network.id`) if err != nil { return nil, errors.New("read network exits") } defer rows.Close() exits := []NetworkExit{} for rows.Next() { exit, err := scanNetworkExit(rows) if err != nil { return nil, err } exits = append(exits, exit) } return exits, rows.Err() } func (s *Store) GetNetworkExit(ctx context.Context, id string) (NetworkExit, error) { if !exitIDPattern.MatchString(id) { return NetworkExit{}, ErrInvalid } return scanNetworkExit(s.db.QueryRowContext(ctx, networkExitSelect+` WHERE network.exit_id = $1`, id)) } func (s *Store) GetNetworkExitAccess(ctx context.Context, id string) (NetworkExitAccess, error) { exit, err := s.GetNetworkExit(ctx, id) return NetworkExitAccess{NetworkExit: exit}, err } const networkExitSelect = ` SELECT network.exit_id, network.protocol, network.host, network.port, network.username, network.password, COALESCE(host(network.expected_public_ip), ''), network.expected_region, COALESCE(host(network.observed_public_ip), ''), network.observed_region, network.health_status, COALESCE(network.last_check_reason, ''), network.version, network.last_checked_at, network.created_at, network.updated_at FROM network_exit network` type rowScanner interface{ Scan(...any) error } func scanNetworkExit(row rowScanner) (NetworkExit, error) { var exit NetworkExit var checked sql.NullTime if err := row.Scan(&exit.ID, &exit.Protocol, &exit.Host, &exit.Port, &exit.Username, &exit.Password, &exit.ExpectedPublicIP, &exit.ExpectedRegion, &exit.ObservedPublicIP, &exit.ObservedRegion, &exit.HealthStatus, &exit.LastCheckReason, &exit.Version, &checked, &exit.CreatedAt, &exit.UpdatedAt); err != nil { return NetworkExit{}, rowError(err) } if checked.Valid { exit.LastCheckedAt = &checked.Time } return exit, nil } // RecordNetworkExitCheck stores only observed identity and a stable reason code. func (s *Store) RecordNetworkExitCheck(ctx context.Context, id string, observation ExitObservation, failureReason string) (NetworkExit, string, error) { if !exitIDPattern.MatchString(id) || !validOptionalRegion(observation.Region) || (observation.PublicIP != "" && net.ParseIP(observation.PublicIP) == nil) || !validExitFailureReason(failureReason) { return NetworkExit{}, "invalid_observation", ErrInvalid } tx, err := s.db.BeginTx(ctx, nil) if err != nil { return NetworkExit{}, "persistence_failed", errors.New("begin network exit check") } defer tx.Rollback() var expectedIP, expectedRegion, oldIP, oldRegion, oldStatus string var version int64 if err := tx.QueryRowContext(ctx, ` SELECT COALESCE(host(expected_public_ip), ''), expected_region, COALESCE(host(observed_public_ip), ''), observed_region, health_status, version FROM network_exit WHERE exit_id = $1 FOR UPDATE`, id). Scan(&expectedIP, &expectedRegion, &oldIP, &oldRegion, &oldStatus, &version); err != nil { return NetworkExit{}, "persistence_failed", rowError(err) } if oldStatus == "disabled" { return NetworkExit{}, "exit_disabled", ErrConflict } reason, status := strings.TrimSpace(failureReason), "unhealthy" if reason == "" && expectedIP != "" && !net.ParseIP(expectedIP).Equal(net.ParseIP(observation.PublicIP)) { reason = "exit_ip_drift" } if reason == "" && expectedRegion != "" && !strings.EqualFold(expectedRegion, observation.Region) { reason = "exit_region_drift" } if reason == "" { reason, status = "exit_healthy", "healthy" } changed := !sameIP(oldIP, observation.PublicIP) || !strings.EqualFold(oldRegion, observation.Region) || oldStatus != status if changed { version++ } if _, err := tx.ExecContext(ctx, ` UPDATE network_exit SET observed_public_ip = NULLIF($2, '')::inet, observed_region = $3, health_status = $4, last_check_reason = $5, version = $6, last_checked_at = now(), updated_at = now() WHERE exit_id = $1`, id, observation.PublicIP, observation.Region, status, reason, version); err != nil { return NetworkExit{}, "persistence_failed", errors.New("record network exit check") } if err := commitHub(tx); err != nil { return NetworkExit{}, "persistence_failed", err } exit, err := s.GetNetworkExit(ctx, id) return exit, reason, err } func validExitFailureReason(reason string) bool { switch reason { case "", "credential_unavailable", "credential_invalid", "proxy_auth_failed", "proxy_check_failed", "exit_observation_invalid": return true default: return false } } func sameIP(left, right string) bool { if left == "" || right == "" { return left == right } return net.ParseIP(left).Equal(net.ParseIP(right)) } func (s *Store) DisableNetworkExit(ctx context.Context, id string) (NetworkExit, error) { if !exitIDPattern.MatchString(id) { return NetworkExit{}, ErrInvalid } tx, err := s.db.BeginTx(ctx, nil) if err != nil { return NetworkExit{}, errors.New("begin network exit disable") } defer tx.Rollback() var oldStatus string if err := tx.QueryRowContext(ctx, `SELECT health_status FROM network_exit WHERE exit_id = $1 FOR UPDATE`, id).Scan(&oldStatus); err != nil { return NetworkExit{}, rowError(err) } var active bool if err := tx.QueryRowContext(ctx, `SELECT EXISTS ( SELECT 1 FROM browser_env b JOIN network_exit n ON n.id = b.exit_id WHERE n.exit_id = $1 AND b.runtime_id IS NOT NULL AND b.runtime_lease_until > now() )`, id).Scan(&active); err != nil { return NetworkExit{}, errors.New("check network exit activity") } if active { return NetworkExit{}, ErrConflict } if oldStatus != "disabled" { if _, err := tx.ExecContext(ctx, ` UPDATE network_exit SET health_status = 'disabled', last_check_reason = 'exit_disabled', version = version + 1, updated_at = now() WHERE exit_id = $1`, id); err != nil { return NetworkExit{}, errors.New("disable network exit") } if _, err := tx.ExecContext(ctx, ` UPDATE social_account account SET status = 'paused', paused_at = COALESCE(paused_at, now()), version = account.version + 1, updated_at = now() FROM browser_env environment JOIN network_exit network ON network.id = environment.exit_id WHERE network.exit_id = $1 AND environment.account_id = account.id`, id); err != nil { return NetworkExit{}, errors.New("invalidate network exit accounts") } } if err := commitHub(tx); err != nil { return NetworkExit{}, err } return s.GetNetworkExit(ctx, id) } func (s *Store) CreateBoundEnv(ctx context.Context, env Env, accountID, exitID string) (EnvironmentContext, bool, error) { env.Alias, env.Name = strings.TrimSpace(env.Alias), strings.TrimSpace(env.Name) if !aliasPattern.MatchString(env.Alias) || !validDisplayName(env.Name) || !aliasPattern.MatchString(accountID) || (exitID != "" && !exitIDPattern.MatchString(exitID)) || !gatewayNamePattern.MatchString(env.Gateway) || env.Fingerprint.ProxyServer != "" { return EnvironmentContext{}, false, ErrInvalid } if err := env.Fingerprint.Validate(); err != nil { return EnvironmentContext{}, false, ErrInvalid } encoded, _ := json.Marshal(env.Fingerprint) tx, err := s.db.BeginTx(ctx, nil) if err != nil { return EnvironmentContext{}, false, errors.New("begin bound environment create") } defer tx.Rollback() var existingAlias, existingExit string err = tx.QueryRowContext(ctx, ` SELECT environment.alias, COALESCE(network.exit_id, '') FROM browser_env environment LEFT JOIN network_exit network ON network.id = environment.exit_id WHERE environment.account_id = (SELECT account.id FROM social_account account WHERE account.account_id = $1) FOR UPDATE OF environment`, accountID). Scan(&existingAlias, &existingExit) if err == nil { if existingAlias != env.Alias || existingExit != exitID { return EnvironmentContext{}, false, ErrConflict } if err := tx.Commit(); err != nil { return EnvironmentContext{}, false, errors.New("commit existing environment lookup") } context, err := s.GetEnvironmentContext(ctx, env.Alias) if err != nil || context.Name != env.Name || context.Gateway != env.Gateway { return EnvironmentContext{}, false, ErrConflict } return context, false, nil } if !errors.Is(err, sql.ErrNoRows) { return EnvironmentContext{}, false, publicDatabaseError(err) } var created string if err := tx.QueryRowContext(ctx, ` INSERT INTO browser_env (alias, name, gateway_id, fingerprint, account_id, exit_id, version) SELECT $1, $2, gateway.id, jsonb_set($4::jsonb, '{seed}', to_jsonb(account.id + 1000)), account.id, network.id, 1 FROM social_account account JOIN gateway ON gateway.name = $3 LEFT JOIN network_exit network ON network.exit_id = NULLIF($6, '') WHERE account.account_id = $5 AND account.status = 'paused' AND ($6 = '' OR network.health_status = 'healthy') RETURNING alias`, env.Alias, env.Name, env.Gateway, encoded, accountID, exitID).Scan(&created); err != nil { return EnvironmentContext{}, false, rowError(err) } if _, err := tx.ExecContext(ctx, `UPDATE social_account SET version = version + 1, updated_at = now() WHERE account_id = $1`, accountID); err != nil { return EnvironmentContext{}, false, errors.New("version bound account") } if err := commitHub(tx); err != nil { return EnvironmentContext{}, false, err } context, err := s.GetEnvironmentContext(ctx, env.Alias) return context, true, err } func (s *Store) GetEnvironmentContext(ctx context.Context, alias string) (EnvironmentContext, error) { if !aliasPattern.MatchString(alias) { return EnvironmentContext{}, ErrInvalid } tx, err := s.db.BeginTx(ctx, nil) if err != nil { return EnvironmentContext{}, errors.New("begin environment context read") } defer tx.Rollback() var result EnvironmentContext var encoded []byte var expectedIP, observedIP string var checked sql.NullTime var runtimeID, runtimeNetworkID, runtimeNodeID, cleanupRuntimeID, cleanupNetworkID sql.NullString var cleanupBindingVersion sql.NullInt64 err = tx.QueryRowContext(ctx, ` SELECT environment.alias, environment.name, gateway.name, environment.fingerprint, environment.created_at, account.account_id, account.status, environment.version, environment.runtime_cleanup_pending, environment.runtime_cleanup_binding_version, environment.runtime_cleanup_runtime_id, environment.runtime_cleanup_network_id, COALESCE(network.exit_id, ''), COALESCE(network.protocol, ''), COALESCE(network.host, ''), COALESCE(network.port, 0), COALESCE(host(network.expected_public_ip), ''), COALESCE(network.expected_region, ''), COALESCE(host(network.observed_public_ip), ''), COALESCE(network.observed_region, ''), COALESCE(network.health_status, 'unchecked'), COALESCE(network.last_check_reason, ''), COALESCE(network.version, 0), network.last_checked_at, COALESCE(network.created_at, to_timestamp(0)), COALESCE(network.updated_at, to_timestamp(0)), COALESCE(environment.runtime_id, ''), COALESCE(environment.runtime_network_id, ''), COALESCE(environment.runtime_node_id, '') FROM browser_env environment JOIN social_account account ON account.id = environment.account_id JOIN gateway ON gateway.id = environment.gateway_id LEFT JOIN network_exit network ON network.id = environment.exit_id WHERE environment.alias = $1`, alias). Scan(&result.Alias, &result.Name, &result.Gateway, &encoded, &result.CreatedAt, &result.AccountID, &result.AccountStatus, &result.BindingVersion, &result.RuntimeCleanupPending, &cleanupBindingVersion, &cleanupRuntimeID, &cleanupNetworkID, &result.Exit.ID, &result.Exit.Protocol, &result.Exit.Host, &result.Exit.Port, &expectedIP, &result.Exit.ExpectedRegion, &observedIP, &result.Exit.ObservedRegion, &result.Exit.HealthStatus, &result.Exit.LastCheckReason, &result.Exit.Version, &checked, &result.Exit.CreatedAt, &result.Exit.UpdatedAt, &runtimeID, &runtimeNetworkID, &runtimeNodeID) if err != nil { return EnvironmentContext{}, rowError(err) } if err := json.Unmarshal(encoded, &result.Fingerprint); err != nil { return EnvironmentContext{}, errors.New("decode bound environment fingerprint") } result.Fingerprint.ProxyServer = "" result.Fingerprint.DisableNonProxiedUDP = false result.Exit.ExpectedPublicIP, result.Exit.ObservedPublicIP = expectedIP, observedIP if checked.Valid { result.Exit.LastCheckedAt = &checked.Time } result.RuntimeID, result.RuntimeNetworkID, result.RuntimeNodeID = runtimeID.String, runtimeNetworkID.String, runtimeNodeID.String if result.RuntimeCleanupPending { result.RuntimeCleanupBindingVersion = cleanupBindingVersion.Int64 result.RuntimeCleanupRuntimeID, result.RuntimeCleanupNetworkID = cleanupRuntimeID.String, cleanupNetworkID.String } if err := commitHub(tx); err != nil { return EnvironmentContext{}, err } return result, nil } func (s *Store) GetEnvironmentContextForAccount(ctx context.Context, accountID string) (EnvironmentContext, error) { if !aliasPattern.MatchString(accountID) { return EnvironmentContext{}, ErrInvalid } var alias string if err := s.db.QueryRowContext(ctx, ` SELECT environment.alias FROM browser_env environment JOIN social_account account ON account.id = environment.account_id WHERE account.account_id = $1`, accountID).Scan(&alias); err != nil { return EnvironmentContext{}, rowError(err) } return s.GetEnvironmentContext(ctx, alias) } // UpdateAccountFingerprint 更新账号绑定环境的指纹参数并返回最新环境上下文。 // seed 保持账号派生值不可改,代理由网络出口管理不可存(直连环境强制空代理)。 func (s *Store) UpdateAccountFingerprint(ctx context.Context, accountID string, fingerprint Fingerprint) (EnvironmentContext, error) { // 门禁对齐 CreateBoundEnv:存储层拒绝携带代理的指纹(代理由网络出口管理,传错入口直接报错)。 if !aliasPattern.MatchString(accountID) || fingerprint.ProxyServer != "" || fingerprint.Validate() != nil { return EnvironmentContext{}, ErrInvalid } tx, err := s.db.BeginTx(ctx, nil) if err != nil { return EnvironmentContext{}, errors.New("begin account fingerprint update") } defer tx.Rollback() var alias string var encoded []byte err = tx.QueryRowContext(ctx, ` SELECT environment.alias, environment.fingerprint FROM browser_env environment JOIN social_account account ON account.id = environment.account_id WHERE account.account_id = $1 FOR UPDATE OF environment`, accountID).Scan(&alias, &encoded) if err != nil { return EnvironmentContext{}, rowError(err) } var stored Fingerprint if err := json.Unmarshal(encoded, &stored); err != nil { return EnvironmentContext{}, errors.New("decode environment fingerprint") } fingerprint.Seed = stored.Seed updated, err := json.Marshal(fingerprint) if err != nil { return EnvironmentContext{}, errors.New("encode environment fingerprint") } if _, err := tx.ExecContext(ctx, `UPDATE browser_env SET fingerprint = $2 WHERE alias = $1`, alias, updated); err != nil { return EnvironmentContext{}, errors.New("update environment fingerprint") } if err := commitHub(tx); err != nil { return EnvironmentContext{}, err } return s.GetEnvironmentContext(ctx, alias) } const runtimeUseLeaseDuration = time.Minute // MissingRuntimeID 标记「创建结果未知」的清理代:网关侧无实物 ID 可供 fence, // 清理时按别名反查网关代际(对外契约:/v1/browsers 的 runtime_id 哨兵值)。 const MissingRuntimeID = "runtime-not-found" func (s *Store) SetRuntimeNode(ctx context.Context, alias, runtimeID, nodeID string) error { if !aliasPattern.MatchString(alias) || !exitIDPattern.MatchString(runtimeID) || !nodeIDPattern.MatchString(nodeID) { return ErrInvalid } result, err := s.db.ExecContext(ctx, ` UPDATE browser_env SET runtime_node_id = $3 WHERE alias = $1 AND runtime_id = $2`, alias, runtimeID, nodeID) if err != nil { return publicDatabaseError(err) } affected, err := result.RowsAffected() if err != nil { return err } if affected != 1 { return ErrConflict } return nil } func releaseExpiredRuntime(ctx context.Context, tx *sql.Tx, alias string) error { var accountID, exitID string var bindingVersion int64 err := tx.QueryRowContext(ctx, ` SELECT account.account_id, COALESCE(network.exit_id, ''), environment.version FROM browser_env environment JOIN social_account account ON account.id = environment.account_id LEFT JOIN network_exit network ON network.id = environment.exit_id WHERE environment.alias = $1 AND environment.runtime_id IS NOT NULL AND environment.runtime_lease_until <= now() FOR UPDATE OF environment`, alias). Scan(&accountID, &exitID, &bindingVersion) if errors.Is(err, sql.ErrNoRows) { return nil } if err != nil { return err } if _, err := tx.ExecContext(ctx, ` UPDATE browser_env SET runtime_id = NULL, runtime_lease_until = NULL, runtime_network_id = NULL WHERE alias = $1`, alias); err != nil { return err } return appendRuntimeAudit(ctx, tx, "runtime_released", accountID, alias, exitID, bindingVersion) } func validateEnvironmentRebind(ctx context.Context, tx *sql.Tx, alias, exitID string, expectedBindingVersion int64) (string, error) { var accountID string var bindingVersion int64 err := tx.QueryRowContext(ctx, ` SELECT account.account_id, environment.version FROM browser_env environment JOIN social_account account ON account.id = environment.account_id WHERE environment.alias = $1 AND account.status = 'paused' AND NOT environment.runtime_cleanup_pending FOR UPDATE OF environment, account`, alias).Scan(&accountID, &bindingVersion) if errors.Is(err, sql.ErrNoRows) { return "", ErrConflict } if err != nil { return "", publicDatabaseError(err) } if bindingVersion != expectedBindingVersion { return "", ErrConflict } if err := releaseExpiredRuntime(ctx, tx, alias); err != nil { return "", errors.New("expire runtime before rebind") } var allowed bool if err := tx.QueryRowContext(ctx, ` SELECT EXISTS (SELECT 1 FROM network_exit WHERE exit_id = $1 AND health_status = 'healthy') AND NOT EXISTS (SELECT 1 FROM browser_env other JOIN social_account other_account ON other_account.id = other.account_id WHERE other_account.account_id = $2 AND other.runtime_id IS NOT NULL)`, exitID, accountID).Scan(&allowed); err != nil { return "", errors.New("check environment rebind") } if !allowed { return "", ErrConflict } return accountID, nil } func (s *Store) ValidateEnvironmentRebind(ctx context.Context, alias, exitID string, expectedBindingVersion int64) error { if !aliasPattern.MatchString(alias) || !exitIDPattern.MatchString(exitID) || expectedBindingVersion < 1 { return ErrInvalid } tx, err := s.db.BeginTx(ctx, nil) if err != nil { return errors.New("begin environment rebind validation") } defer tx.Rollback() if _, err := validateEnvironmentRebind(ctx, tx, alias, exitID, expectedBindingVersion); err != nil { return err } if err := commitHub(tx); err != nil { return err } return nil } func (s *Store) RebindEnvironment(ctx context.Context, alias, exitID, runtimeID string, expectedBindingVersion int64, networkIDs ...string) (EnvironmentContext, error) { networkID := "" if len(networkIDs) == 1 { networkID = networkIDs[0] } if !aliasPattern.MatchString(alias) || !exitIDPattern.MatchString(exitID) || (runtimeID != "" && !exitIDPattern.MatchString(runtimeID)) || (networkID != "" && !exitIDPattern.MatchString(networkID)) || len(networkIDs) > 1 || expectedBindingVersion < 1 { return EnvironmentContext{}, ErrInvalid } tx, err := s.db.BeginTx(ctx, nil) if err != nil { return EnvironmentContext{}, errors.New("begin environment rebind") } defer tx.Rollback() accountID, err := validateEnvironmentRebind(ctx, tx, alias, exitID, expectedBindingVersion) if err != nil { return EnvironmentContext{}, err } if _, err := tx.ExecContext(ctx, `UPDATE browser_env SET exit_id = (SELECT network.id FROM network_exit network WHERE network.exit_id = $2), version = version + 1, updated_at = now() WHERE alias = $1`, alias, exitID); err != nil { return EnvironmentContext{}, errors.New("update environment binding") } if _, err := tx.ExecContext(ctx, `UPDATE social_account SET version = version + 1, updated_at = now() WHERE account_id = $1`, accountID); err != nil { return EnvironmentContext{}, errors.New("version rebound account") } if runtimeID != "" { if _, err := tx.ExecContext(ctx, ` UPDATE browser_env SET runtime_id = $2, runtime_lease_until = now() + interval '1 minute', runtime_network_id = NULLIF($3, '') WHERE alias = $1`, alias, runtimeID, networkID); err != nil { return EnvironmentContext{}, publicDatabaseError(err) } if err := appendRuntimeAudit(ctx, tx, "runtime_bound", accountID, alias, exitID, expectedBindingVersion+1); err != nil { return EnvironmentContext{}, err } } if err := commitHub(tx); err != nil { return EnvironmentContext{}, err } return s.GetEnvironmentContext(ctx, alias) } func (s *Store) ActivateRuntime(ctx context.Context, alias, runtimeID string, bindingVersion int64, exitID string, networkIDs ...string) (EnvironmentContext, error) { networkID := "" if len(networkIDs) == 1 { networkID = networkIDs[0] } if !aliasPattern.MatchString(alias) || !exitIDPattern.MatchString(runtimeID) || bindingVersion < 1 || (exitID != "" && !exitIDPattern.MatchString(exitID)) || !exitIDPattern.MatchString(networkID) || len(networkIDs) != 1 { return EnvironmentContext{}, ErrInvalid } tx, err := s.db.BeginTx(ctx, nil) if err != nil { return EnvironmentContext{}, errors.New("begin runtime activation") } defer tx.Rollback() var accountID, currentExitID, accountStatus string var currentBindingVersion int64 var cleanupPending bool err = tx.QueryRowContext(ctx, ` SELECT account.account_id, environment.version, COALESCE(network.exit_id, ''), environment.runtime_cleanup_pending, account.status FROM browser_env environment JOIN social_account account ON account.id = environment.account_id LEFT JOIN network_exit network ON network.id = environment.exit_id WHERE environment.alias = $1 FOR UPDATE OF environment, account`, alias). Scan(&accountID, ¤tBindingVersion, ¤tExitID, &cleanupPending, &accountStatus) if err != nil { return EnvironmentContext{}, rowError(err) } if cleanupPending || accountStatus != "active" || currentBindingVersion != bindingVersion || currentExitID != exitID { return EnvironmentContext{}, ErrConflict } if err := releaseExpiredRuntime(ctx, tx, alias); err != nil { return EnvironmentContext{}, errors.New("expire runtime before activation") } var existingRuntimeID, existingNetworkID string if err := tx.QueryRowContext(ctx, ` SELECT COALESCE(runtime_id, ''), COALESCE(runtime_network_id, '') FROM browser_env WHERE alias = $1 AND runtime_id IS NOT NULL`, alias).Scan(&existingRuntimeID, &existingNetworkID); err != nil && !errors.Is(err, sql.ErrNoRows) { return EnvironmentContext{}, publicDatabaseError(err) } if existingRuntimeID != "" && (existingRuntimeID != runtimeID || existingNetworkID != networkID) { return EnvironmentContext{}, ErrConflict } if existingRuntimeID == "" { if _, err := tx.ExecContext(ctx, ` UPDATE browser_env SET runtime_id = $2, runtime_lease_until = now() + interval '1 minute', runtime_network_id = NULLIF($3, '') WHERE alias = $1`, alias, runtimeID, networkID); err != nil { return EnvironmentContext{}, publicDatabaseError(err) } if err := appendRuntimeAudit(ctx, tx, "runtime_bound", accountID, alias, exitID, bindingVersion); err != nil { return EnvironmentContext{}, err } } else if _, err := tx.ExecContext(ctx, `UPDATE browser_env SET runtime_lease_until = now() + interval '1 minute' WHERE alias = $1`, alias); err != nil { return EnvironmentContext{}, errors.New("renew environment runtime") } if err := commitHub(tx); err != nil { return EnvironmentContext{}, err } return s.GetEnvironmentContext(ctx, alias) } func (s *Store) ReleaseRuntime(ctx context.Context, environment EnvironmentContext) error { if !aliasPattern.MatchString(environment.Alias) || environment.BindingVersion < 1 || (environment.RuntimeID != "" && !exitIDPattern.MatchString(environment.RuntimeID)) { return ErrInvalid } if environment.RuntimeID == "" { return nil } tx, err := s.db.BeginTx(ctx, nil) if err != nil { return errors.New("begin runtime release") } defer tx.Rollback() var accountID, exitID string var bindingVersion int64 err = tx.QueryRowContext(ctx, ` SELECT account.account_id, COALESCE(network.exit_id, ''), environment.version FROM browser_env environment JOIN social_account account ON account.id = environment.account_id LEFT JOIN network_exit network ON network.id = environment.exit_id WHERE environment.alias = $1 AND environment.version = $2 AND environment.runtime_id = $3 FOR UPDATE OF environment, account`, environment.Alias, environment.BindingVersion, environment.RuntimeID). Scan(&accountID, &exitID, &bindingVersion) if errors.Is(err, sql.ErrNoRows) { return ErrConflict } if err != nil { return errors.New("release environment runtime") } if _, err := tx.ExecContext(ctx, ` UPDATE browser_env SET runtime_id = NULL, runtime_lease_until = NULL, runtime_network_id = NULL WHERE alias = $1`, environment.Alias); err != nil { return errors.New("release environment runtime") } if err := appendRuntimeAudit(ctx, tx, "runtime_released", accountID, environment.Alias, exitID, bindingVersion); err != nil { return err } return commitHub(tx) } func (s *Store) SetRuntimeCleanupPending(ctx context.Context, environment EnvironmentContext, pending bool) error { if !aliasPattern.MatchString(environment.Alias) || environment.BindingVersion < 1 || environment.RuntimeCleanupBindingVersion < 1 || (pending && environment.RuntimeCleanupRuntimeID == "") || (environment.RuntimeCleanupRuntimeID != "" && !exitIDPattern.MatchString(environment.RuntimeCleanupRuntimeID)) || (environment.RuntimeCleanupNetworkID != "" && !exitIDPattern.MatchString(environment.RuntimeCleanupNetworkID)) { return ErrInvalid } tx, err := s.db.BeginTx(ctx, nil) if err != nil { return errors.New("begin runtime cleanup state update") } defer tx.Rollback() var currentPending bool var accountID string var exitID sql.NullString var cleanupBindingVersion sql.NullInt64 var cleanupRuntimeID, cleanupNetworkID sql.NullString if err := tx.QueryRowContext(ctx, ` SELECT environment.runtime_cleanup_pending, environment.runtime_cleanup_binding_version, environment.runtime_cleanup_runtime_id, environment.runtime_cleanup_network_id, account.account_id, COALESCE(network.exit_id, '') FROM browser_env environment JOIN social_account account ON account.id = environment.account_id LEFT JOIN network_exit network ON network.id = environment.exit_id WHERE environment.alias = $1 AND environment.version = $2 FOR UPDATE OF environment, account`, environment.Alias, environment.BindingVersion). Scan(¤tPending, &cleanupBindingVersion, &cleanupRuntimeID, &cleanupNetworkID, &accountID, &exitID); errors.Is(err, sql.ErrNoRows) { return ErrConflict } else if err != nil { return publicDatabaseError(err) } if currentPending { if cleanupBindingVersion.Int64 != environment.RuntimeCleanupBindingVersion || cleanupRuntimeID.String != environment.RuntimeCleanupRuntimeID || cleanupNetworkID.String != environment.RuntimeCleanupNetworkID { return ErrConflict } if pending { return commitHub(tx) } } else if !pending { return commitHub(tx) } if pending && !currentPending { var runtimeID string err := tx.QueryRowContext(ctx, ` SELECT COALESCE(runtime_id, '') FROM browser_env WHERE alias = $1 AND version = $2`, environment.Alias, environment.BindingVersion).Scan(&runtimeID) if err != nil { return publicDatabaseError(err) } if runtimeID != environment.RuntimeCleanupRuntimeID && environment.RuntimeCleanupRuntimeID != MissingRuntimeID { // fence 记录清理目标代:目标=当前活跃 runtime(显式停止)、未知代(创建结果未知,登记时释放活跃实例) // 或无活跃实例(创建失败)。其余视为过期代际请求。 if runtimeID != "" { return ErrConflict } } if runtimeID != "" { result, err := tx.ExecContext(ctx, ` UPDATE browser_env SET runtime_id = NULL, runtime_lease_until = NULL, runtime_network_id = NULL WHERE alias = $1 AND version = $2 AND runtime_id = $3`, environment.Alias, environment.BindingVersion, runtimeID) if err != nil { return errors.New("release runtime for pending cleanup") } if affected, err := result.RowsAffected(); err != nil || affected != 1 { return ErrConflict } if err := appendRuntimeAudit(ctx, tx, "runtime_released", accountID, environment.Alias, exitID.String, environment.BindingVersion); err != nil { return err } } } if pending { _, err = tx.ExecContext(ctx, ` UPDATE browser_env SET runtime_cleanup_pending = true, runtime_cleanup_binding_version = $2, runtime_cleanup_runtime_id = NULLIF($3, ''), runtime_cleanup_network_id = NULLIF($4, ''), updated_at = now() WHERE alias = $1`, environment.Alias, environment.RuntimeCleanupBindingVersion, environment.RuntimeCleanupRuntimeID, environment.RuntimeCleanupNetworkID) } else { _, err = tx.ExecContext(ctx, ` UPDATE browser_env SET runtime_cleanup_pending = false, runtime_cleanup_binding_version = NULL, runtime_cleanup_runtime_id = NULL, runtime_cleanup_network_id = NULL, updated_at = now() WHERE alias = $1`, environment.Alias) } if err != nil { return errors.New("update runtime cleanup state") } return commitHub(tx) } func appendRuntimeAudit(ctx context.Context, tx *sql.Tx, eventType, accountID, alias, exitID string, bindingVersion int64) error { _, err := tx.ExecContext(ctx, ` INSERT INTO audit_event (event_type, account_id, browser_env_id, exit_id, binding_version, actor, reason_code) VALUES ($1, (SELECT account.id FROM social_account account WHERE account.account_id = $2), (SELECT environment.id FROM browser_env environment WHERE environment.alias = $3), (SELECT network.id FROM network_exit network WHERE network.exit_id = NULLIF($4, '')), $5, 'local-user', $1)`, eventType, accountID, alias, exitID, bindingVersion) if err != nil { return errors.New("append runtime audit event") } return nil } func (s *Store) AppendEnvironmentAction(ctx context.Context, eventType string, action EnvironmentAction) error { if (eventType != "environment_action_requested" && eventType != "environment_action_finished") || !exitIDPattern.MatchString(action.OperationID) || action.Action == "" || action.ReasonCode == "" || (eventType == "environment_action_finished" && action.Outcome != "succeeded" && action.Outcome != "failed" && action.Outcome != "unknown") { return ErrInvalid } _, err := s.db.ExecContext(ctx, ` INSERT INTO audit_event (event_type, account_id, browser_env_id, exit_id, binding_version, actor, reason_code, operation_id, action, outcome) VALUES ($1, (SELECT account.id FROM social_account account WHERE account.account_id = NULLIF($2, '')), (SELECT environment.id FROM browser_env environment WHERE environment.alias = NULLIF($3, '')), (SELECT network.id FROM network_exit network WHERE network.exit_id = NULLIF($4, '')), NULLIF($5, 0), 'local-user', $6, $7, $8, NULLIF($9, ''))`, eventType, action.AccountID, action.BrowserEnvAlias, action.NetworkExitID, action.BindingVersion, action.ReasonCode, action.OperationID, action.Action, action.Outcome) if err != nil { return errors.New("append environment action") } return nil } func NewOperationID() string { return "operation-" + newHubID() } func newHubID() string { var value [12]byte _, _ = rand.Read(value[:]) return hex.EncodeToString(value[:]) }