using System.Net.Http;
using System.Text.Json.Nodes;
using Dpz.Core.Public.Entity.Steam;
using Dpz.Core.Public.ViewModel.Response.Steam;
using Microsoft.AspNetCore.WebUtilities;

namespace Dpz.Core.Service.RepositoryServiceImpl;

public class SteamGameService(
    IRepository<SteamGame> repository,
    IMapper mapper,
    HttpClient httpClient,
    IConfiguration configuration,
    ILogger<SteamGameService> logger,
    IFusionCache fusionCache
) : ISteamGameService
{
    private readonly Lazy<string> _lazyCdnHost = new(() =>
    {
        var cdnHost = configuration["S3:CdnHost"];
        return string.IsNullOrWhiteSpace(cdnHost)
            ? throw new InvalidConfigurationException("S3:CdnHost is empty or null")
            : cdnHost;
    });

    private readonly Lazy<string> _lazyLogoUrlTemplate = new(() =>
        "https://shared.akamai.steamstatic.com/store_item_assets/steam/apps/{0}/header.jpg"
    );

    private readonly Lazy<string> _lazySteamId = new(() =>
    {
        var steamId = configuration["Steam:SteamId"];
        return string.IsNullOrWhiteSpace(steamId)
            ? throw new InvalidConfigurationException("Steam:SteamId is empty or null")
            : steamId;
    });

    private readonly Lazy<string> _lazySteamAppKey = new(() =>
    {
        var appKey = configuration["Steam:AppKey"];
        return string.IsNullOrWhiteSpace(appKey)
            ? throw new InvalidConfigurationException("Steam:SteamAppKey is empty or null")
            : appKey;
    });

    private readonly Lazy<List<int>> _filterIgnoreGameIds = new(() =>
        configuration.GetSection("Steam:IgnoreFilterIds").Get<List<int>>() ?? []
    );

    // API请求延迟时间,防止请求过于频繁
    private const int ApiDelayMs = 500;

    /// <summary>
    /// 当前实例ID,用于区分不同的实例
    /// </summary>
    private readonly Lazy<string> _instanceId = new(() => Guid.NewGuid().ToString("D"));

    /// <summary>
    /// 当前实例ID的缓存键
    /// </summary>
    private const string RunningInstanceIdKey = "SteamGameService.Running.InstanceId";

    public event LogoDownload? OnLogoDownloadComplete;
    public event AchievementIconGrayDownload? OnAchievementIconGrayDownloadComplete;
    public event AchievementIconDownload? OnAchievementIconDownloadComplete;

    [InvalidateCache(Methods = [nameof(GetGamesAsync), nameof(GetGameAsync)])]
    public async Task UpdateGamesAsync(CancellationToken cancellationToken = default)
    {
        var runningInstanceId = await fusionCache.TryGetAsync<string>(
            RunningInstanceIdKey,
            token: cancellationToken
        );
        if (runningInstanceId.HasValue && runningInstanceId.Value != _instanceId.Value)
        {
            logger.LogInformation(
                "当前已有实例正在运行,跳过此次更新,实例ID:{RunningInstanceId},当前实例ID:{CurrentInstanceId}",
                runningInstanceId.Value,
                _instanceId.Value
            );
            return;
        }

        await fusionCache.SetAsync(
            RunningInstanceIdKey,
            _instanceId.Value,
            TimeSpan.FromHours(3),
            token: cancellationToken
        );

        try
        {
            logger.LogInformation("开始更新Steam游戏数据,实例ID:{InstanceId}", _instanceId.Value);
            var dbGames = await repository.SearchFor(x => true).ToListAsync(cancellationToken);
            var games = await FetchGamesAsync(dbGames, cancellationToken);

            // 安全检查:如果获取到空列表,直接返回,避免误删数据
            if (games.Count == 0)
            {
                logger.LogWarning("获取到的Steam游戏数量为0,跳过此次更新以防止误删数据");
                return;
            }

            // 新增游戏
            var newGames = games.ExceptBy(dbGames.Select(x => x.Id), x => x.Id).ToList();
            if (newGames.Count > 0)
            {
                logger.LogInformation("发现{Count}个新游戏,准备添加", newGames.Count);
                await repository.InsertAsync(newGames, cancellationToken);
                // var batchSize =
                //     newGames.Count % 10 > 0 ? newGames.Count / 10 + 1 : newGames.Count / 10;
                //
                // for (var i = 0; i < batchSize; i++)
                // {
                //     var startIndex = i * 10;
                //     var count = Math.Min(10, newGames.Count - startIndex);
                //     var batch = newGames.GetRange(startIndex, count);
                //     logger.LogInformation(
                //         "批量添加游戏,开始索引:{StartIndex},数量:{Count}",
                //         startIndex,
                //         count
                //     );
                //     await repository.InsertAsync(batch);
                //     await Task.Delay(1000);
                // }
            }

            // DB已有游戏,待更新
            var updateGames = games.IntersectBy(dbGames.Select(x => x.Id), x => x.Id).ToList();
            if (updateGames.Count > 0)
            {
                logger.LogInformation("开始更新{Count}个已有游戏", updateGames.Count);
                await Parallel.ForEachAsync(
                    updateGames,
                    cancellationToken,
                    async (game, _) => await UpdateGameAsync(dbGames, game)
                );
            }

            // 删除已删除的游戏
            var deletedGames = dbGames.ExceptBy(games.Select(x => x.Id), x => x.Id).ToList();
            if (deletedGames.Count > 0)
            {
                logger.LogInformation("发现{Count}个已删除游戏,准备删除", deletedGames.Count);

                var filter = Builders<SteamGame>.Filter.In(
                    x => x.Id,
                    deletedGames.Select(x => x.Id)
                );
                await repository.DeleteAsync(filter, cancellationToken);
            }

            logger.LogInformation(
                "Steam游戏数据更新完成,新增:{NewCount},更新:{UpdateCount},删除:{DeletedCount},实例ID:{InstanceId}",
                newGames.Count,
                updateGames.Count,
                deletedGames.Count,
                _instanceId.Value
            );
        }
        catch (Exception ex)
        {
            logger.LogError(
                ex,
                "更新Steam游戏数据时发生异常,实例ID:{InstanceId}",
                _instanceId.Value
            );
            throw;
        }
        finally
        {
            // 确保在任何情况下都释放锁
            try
            {
                await fusionCache.RemoveAsync(RunningInstanceIdKey, token: cancellationToken);
                logger.LogDebug("已释放Steam游戏更新锁,实例ID:{InstanceId}", _instanceId.Value);
            }
            catch (Exception ex)
            {
                logger.LogError(
                    ex,
                    "释放Steam游戏更新锁时发生异常,实例ID:{InstanceId}",
                    _instanceId.Value
                );
            }
        }
    }

    private async Task UpdateGameAsync(List<SteamGame> dbGames, SteamGame updateGame)
    {
        var dbGame = dbGames.Find(x => x.Id == updateGame.Id);
        if (dbGame != null)
        {
            var dic = dbGame.UpdateContent(updateGame);
            if (dic.Count > 0)
            {
                logger.LogInformation(
                    "更新游戏 {GameId}/{GameName} 的数据,更新细节:{@UpdateContent}",
                    updateGame.Id,
                    updateGame.Name,
                    dic
                );
            }
        }
        await repository.UpdateAsync(updateGame);
    }

    private async Task<List<AchievementDetail>> FetchAchievementsAsync(
        int id,
        List<SteamGame> dbGames,
        CancellationToken cancellationToken = default
    )
    {
        var queryString = new Dictionary<string, string?>
        {
            { "steamid", _lazySteamId.Value },
            { "key", _lazySteamAppKey.Value },
            { "format", "json" },
            { "appid", id.ToString() },
            { "l", "schinese" },
        };

        // 添加API请求延迟
        await Task.Delay(ApiDelayMs, cancellationToken);

        logger.LogDebug("获取游戏 {GameId} 的成就数据", id);
        var achievements =
            await ExecuteHttpRequestAsync<List<AchievementDetail>>(
                "/ISteamUserStats/GetSchemaForGame/v2/",
                queryString,
                root => root?["game"]?["availableGameStats"]?["achievements"]
            ) ?? [];

        var unlockAchievements = await FetchUnlockAchievementAsync(queryString);
#if DEBUG
        var options = new ParallelOptions { MaxDegreeOfParallelism = 1 };
#endif
        await Parallel.ForEachAsync(
            achievements,
#if DEBUG
            options,
#endif
            async (achievement, _) =>
            {
                if (
                    achievement.Name != null
                    && unlockAchievements.TryGetValue(achievement.Name, out var unlockTime)
                    && unlockTime > 0
                )
                {
                    achievement.UnlockTime = unlockTime.ToDateTime();
                }

                await DownloadAchievementAsync(dbGames, id, achievement);
                await DownloadAchievementGrayAsync(dbGames, id, achievement);
            }
        );

        logger.LogDebug("游戏 {GameId} 共获取到 {Count} 个成就", id, achievements.Count);
        return achievements;
    }

    private async Task<Dictionary<string, long>> FetchUnlockAchievementAsync(
        Dictionary<string, string?> queryString
    )
    {
        logger.LogDebug("获取已解锁的成就数据");

        // 添加API请求延迟
        await Task.Delay(ApiDelayMs);

        return await ExecuteHttpRequestAsync<Dictionary<string, long>>(
                "/ISteamUserStats/GetPlayerAchievements/v0001/",
                queryString,
                root =>
                {
                    var achievements = root
                        ?["playerstats"]?["achievements"]?.AsArray()
                        .Select(x =>
                        {
                            var name = x?["apiname"]?.GetValue<string>();
                            if (name == null)
                            {
                                logger.LogError(
                                    "Achievement name is null for {SteamId}",
                                    _lazySteamId.Value
                                );
                                return (KeyValuePair<string, long>?)null;
                            }
                            var unlockTime = x?["unlocktime"]?.GetValue<long>() ?? 0;
                            return x?["achieved"]?.GetValue<int>() == 1
                                ? KeyValuePair.Create(name, unlockTime)
                                : KeyValuePair.Create(name, 0L);
                        })
                        .Where(x => x != null)
                        .Select(x => x!.Value)
                        .ToDictionary(x => x.Key, x => x.Value);
                    return JsonValue.Create(achievements);
                }
            ) ?? [];
    }

    private async Task<List<SteamGame>> FetchGamesAsync(
        List<SteamGame> dbGames,
        CancellationToken cancellationToken = default
    )
    {
        logger.LogInformation("开始获取Steam游戏列表");

        var games = await ApplicationTools.RetryAsync(
            async () =>
            {
                var queryString = new Dictionary<string, string?>
                {
                    { "steamid", _lazySteamId.Value },
                    { "key", _lazySteamAppKey.Value },
                    { "include_appinfo", "1" },
                    { "format", "json" },
                };

                // 添加API请求延迟
                await Task.Delay(ApiDelayMs, cancellationToken);

                var result = await ExecuteHttpRequestAsync<List<SteamGame>>(
                    "/IPlayerService/GetOwnedGames/v0001/",
                    queryString,
                    root => root?["response"]?["games"]?.AsArray()
                );

                // 如果返回null,抛出异常以触发重试
                if (result == null)
                {
                    throw new InvalidOperationException("Steam API返回null,触发重试");
                }

                return result
                    // 已失效的游戏,导致后续请求图标、成就失败
                    // 不知道为什么steam不移除它
                    // 屏蔽
                    //.Where(x => x.Id != 692850)
                    .Where(x => !_filterIgnoreGameIds.Value.Contains(x.Id))
                    .ToList();
            },
            retryInterval: TimeSpan.FromMilliseconds(ApiDelayMs),
            maxAttemptCount: 3,
            specificExceptionTypes: [typeof(JsonException)] // JSON错误不重试
        );

        logger.LogInformation("获取到 {Count} 个Steam游戏", games.Count);

#if DEBUG
        var options = new ParallelOptions { MaxDegreeOfParallelism = 1 };
#endif
        await Parallel.ForEachAsync(
            games,
#if DEBUG
            options,
#endif
            async (game, ct) => await UpdateGame(dbGames, game, ct)
        );

        return games;
    }

    private async Task UpdateGame(
        List<SteamGame> dbGames,
        SteamGame game,
        CancellationToken cancellationToken = default
    )
    {
        logger.LogDebug("更新游戏 {GameId}/{GameName} 的基本信息", game.Id, game.Name);
        var achievements = await FetchAchievementsAsync(game.Id, dbGames, cancellationToken);

        var originGame = dbGames.Find(x => x.Id == game.Id);

        if (achievements is { Count: > 0 })
        {
            game.Achievements = achievements;
            game.LastUpdateTime = DateTime.Now;
        }
        else if (originGame != null)
        {
            game.Achievements = originGame.Achievements;
        }
        else
        {
            game.Achievements = [];
            game.LastUpdateTime = DateTime.Now;
        }

        if (!AreAssetsDownloaded(originGame))
        {
            logger.LogDebug("下载游戏 {GameId} 的Logo", game.Id);

            // 添加API请求延迟
            await Task.Delay(ApiDelayMs, cancellationToken);

            await using var stream = await HttpDownloadAsync(
                string.Format(_lazyLogoUrlTemplate.Value, game.Id)
            );
            if (stream != null && OnLogoDownloadComplete != null)
            {
                var url = await OnLogoDownloadComplete(stream, game.Id);
                if (!string.IsNullOrWhiteSpace(url))
                {
                    game.ImageLogo = url;
                    game.ImageIcon = url;
                    game.LastUpdateTime = DateTime.Now;
                }
            }
        }
        else
        {
            game.ImageLogo = originGame?.ImageLogo;
            game.ImageIcon = originGame?.ImageIcon;
            logger.LogDebug("游戏 {GameId} Logo 无需更新", game.Id);
        }
    }

    private async Task<Stream?> HttpDownloadAsync(string uri)
    {
        try
        {
            var response = await httpClient.GetAsync(uri);
            response.EnsureSuccessStatusCode();
            var memoryStream = new MemoryStream();
            await response.Content.CopyToAsync(memoryStream);
            memoryStream.Position = 0;
            return memoryStream;
        }
        catch (Exception e)
        {
            logger.LogError(e, "Download failed for {Uri}", uri);
            return null;
        }
    }

    private bool AreAssetsDownloaded(SteamGame? originGame)
    {
        return originGame != null
            && originGame.ImageIcon?.StartsWith(
                _lazyCdnHost.Value,
                StringComparison.OrdinalIgnoreCase
            ) == true
            && originGame.ImageLogo?.StartsWith(
                _lazyCdnHost.Value,
                StringComparison.OrdinalIgnoreCase
            ) == true;
    }

    private async Task DownloadAchievementAsync(
        List<SteamGame> dbGames,
        int id,
        AchievementDetail achievement
    )
    {
        var dbAchievement = dbGames
            .FirstOrDefault(x => x.Id == id)
            ?.Achievements.Find(x => x.Name == achievement.Name);

        if (
            dbAchievement?.Icon?.StartsWith(_lazyCdnHost.Value, StringComparison.OrdinalIgnoreCase)
            == true
        )
        {
            achievement.Icon = dbAchievement.Icon;
            return;
        }

        if (
            string.IsNullOrWhiteSpace(achievement.Icon)
            || string.IsNullOrWhiteSpace(achievement.Name)
        )
        {
            return;
        }

        logger.LogDebug("下载游戏 {GameId} 成就 {AchievementName} 的图标", id, achievement.Name);

        // 添加API请求延迟
        await Task.Delay(ApiDelayMs);

        await using var icon = await HttpDownloadAsync(achievement.Icon);
        var handler = OnAchievementIconDownloadComplete;
        if (icon != null && handler != null)
        {
            var url = await handler.Invoke(icon, id, achievement.Name);
            if (!string.IsNullOrWhiteSpace(url))
            {
                achievement.Icon = url;
            }
        }
    }

    private async Task DownloadAchievementGrayAsync(
        IEnumerable<SteamGame> dbGames,
        int id,
        AchievementDetail achievement
    )
    {
        var dbAchievement = dbGames
            .FirstOrDefault(x => x.Id == id)
            ?.Achievements.Find(x => x.Name == achievement.Name);

        if (
            dbAchievement?.IconGray?.StartsWith(
                _lazyCdnHost.Value,
                StringComparison.OrdinalIgnoreCase
            ) == true
        )
        {
            achievement.IconGray = dbAchievement.IconGray;
            return;
        }
        if (
            string.IsNullOrWhiteSpace(achievement.IconGray)
            || string.IsNullOrWhiteSpace(achievement.Name)
        )
        {
            return;
        }

        logger.LogDebug(
            "下载游戏 {GameId} 成就 {AchievementName} 的灰色图标",
            id,
            achievement.Name
        );

        // 添加API请求延迟
        await Task.Delay(ApiDelayMs);

        await using var icon = await HttpDownloadAsync(achievement.IconGray);
        var handler = OnAchievementIconGrayDownloadComplete;
        if (icon != null && handler != null)
        {
            var url = await handler.Invoke(icon, id, achievement.Name);
            if (!string.IsNullOrWhiteSpace(url))
            {
                achievement.IconGray = url;
            }
        }
    }

    [Cache(ExpirationSeconds = 1 * ExpirationTime.Day)]
    public async Task<List<SteamGameResponse>> GetGamesAsync(
        CancellationToken cancellationToken = default
    )
    {
        var list = await repository.SearchFor(x => true).ToListAsync(cancellationToken);
        var result = mapper.Map<List<SteamGameResponse>>(list);
        return result;
    }

    [Cache(ExpirationSeconds = 1 * ExpirationTime.Day)]
    public async Task<SteamGameResponse?> GetGameAsync(
        int id,
        CancellationToken cancellationToken = default
    )
    {
        var game = await repository
            .SearchFor(x => x.Id == id)
            .FirstOrDefaultAsync(cancellationToken);
        if (game == null)
        {
            logger.LogWarning("未找到游戏 {GameId} 的信息", id);
            return null;
        }
        var result = mapper.Map<SteamGameResponse>(game);
        return result;
    }

    [InvalidateCache(Methods = [nameof(GetGamesAsync), nameof(GetGameAsync)])]
    public Task ClearCacheAsync(CancellationToken cancellationToken = default)
    {
        return Task.CompletedTask;
    }

    private async Task<T?> ExecuteHttpRequestAsync<T>(
        string endpoint,
        Dictionary<string, string?> queryParams,
        Func<JsonNode?, JsonNode?>? dataSelector = null
    )
    {
        var requestUri = QueryHelpers.AddQueryString(endpoint, queryParams);
        var jsonContent = "";
        try
        {
            var response = await httpClient.GetAsync(requestUri);
            response.EnsureSuccessStatusCode();
            jsonContent = await response.Content.ReadAsStringAsync();

            if (string.IsNullOrWhiteSpace(jsonContent))
            {
                logger.LogWarning("Steam API返回空内容,请求地址:{RequestUri}", requestUri);
                return default;
            }

            var jsonNode = JsonNode.Parse(jsonContent);
            if (jsonNode == null)
            {
                logger.LogWarning(
                    "Steam API返回的JSON无法解析,请求地址:{RequestUri},内容:{Content}",
                    requestUri,
                    jsonContent
                );
                return default;
            }

            var data = dataSelector?.Invoke(jsonNode) ?? jsonNode;

            var result = data.Deserialize<T>();
            if (result == null)
            {
                logger.LogWarning(
                    "Steam API数据反序列化失败,请求地址:{RequestUri},数据:{Data}",
                    requestUri,
                    data.ToString()
                );
            }

            return result;
        }
        catch (JsonException e)
        {
            logger.LogError(
                e,
                "Steam API数据解析失败,请求地址:{RequestUri},内容:{Content}",
                requestUri,
                jsonContent
            );
            return default;
        }
        catch (Exception e)
        {
            logger.LogError(e, "Steam API请求失败,请求地址:{RequestUri}", requestUri);
            return default;
        }
    }
}
⚠⚠    以下内容为AI分析的结果,请根据实际情况进行判断。

代码解释

这是一个用于管理和同步 Steam 游戏数据的服务类 SteamGameService,主要功能包括从 Steam API 获取游戏信息、成就数据,以及下载和管理游戏相关图片资源。

核心功能模块

1. 依赖注入与配置初始化

public SteamGameService(
    IRepository<SteamGame> repository,
    IMapper mapper,
    HttpClient httpClient,
    IConfiguration configuration,
    ILogger<SteamGameService> logger,
    IFusionCache fusionCache
)
  • 使用主构造函数注入依赖项
  • 通过 Lazy<T> 延迟加载配置项(CDN 主机、Steam ID、API Key 等)
  • 配置项为空时抛出 InvalidConfigurationException

2. 游戏数据更新(核心方法)

UpdateGamesAsync 方法实现了完整的数据同步流程:

分布式锁机制

var runningInstanceId = await fusionCache.TryGetAsync<string>(RunningInstanceIdKey, ...);
  • 使用 FusionCache 实现分布式锁,防止多实例并发执行
  • 每个实例生成唯一 ID,锁有效期 3 小时
  • finally 块确保锁一定被释放

数据同步逻辑

  1. 新增游戏:从 API 获取的游戏列表中排除数据库已有的
  2. 更新游戏:对数据库中已存在的游戏进行增量更新
  3. 删除游戏:移除 API 中不存在但数据库有的记录
  4. 安全检查:API 返回空列表时跳过更新,防止误删

3. Steam API 交互

API 请求频率控制

private const int ApiDelayMs = 500;
await Task.Delay(ApiDelayMs, cancellationToken);
  • 每次请求前延迟 500ms,避免触发 Steam API 限流

重试机制

var games = await ApplicationTools.RetryAsync(
    async () => { ... },
    retryInterval: TimeSpan.FromMilliseconds(ApiDelayMs),
    maxAttemptCount: 3,
    specificExceptionTypes: [typeof(JsonException)]
);
  • 最多重试 3 次
  • JSON 解析错误不重试(数据格式问题)

通用 HTTP 请求封装

ExecuteHttpRequestAsync<T> 方法:

  • 支持自定义 JSON 数据提取器 (dataSelector)
  • 统一异常处理和日志记录
  • 返回泛型结果,增强复用性

4. 成就数据处理

获取成就详情

private async Task<List<AchievementDetail>> FetchAchievementsAsync(...)
  • /ISteamUserStats/GetSchemaForGame/v2/ 获取成就列表
  • /ISteamUserStats/GetPlayerAchievements/v0001/ 获取解锁状态
  • 合并两个接口数据,填充 UnlockTime 字段

并行下载成就图标

await Parallel.ForEachAsync(achievements, async (achievement, _) => {
    await DownloadAchievementAsync(...);
    await DownloadAchievementGrayAsync(...);
});
  • 同时下载彩色和灰色图标
  • Debug 模式下限制并行度为 1(方便调试)

5. 资源下载与事件通知

委托事件机制

public event LogoDownload? OnLogoDownloadComplete;
public event AchievementIconGrayDownload? OnAchievementIconGrayDownloadComplete;
public event AchievementIconDownload? OnAchievementIconDownloadComplete;
  • 下载完成后触发事件,由外部处理文件上传到 S3
  • 返回 CDN URL 后更新到实体对象

下载优化

private bool AreAssetsDownloaded(SteamGame? originGame)
{
    return originGame?.ImageIcon?.StartsWith(_lazyCdnHost.Value, ...) == true;
}
  • 检查资源是否已存在于 CDN,避免重复下载

6. 缓存策略

方法级缓存

[Cache(ExpirationSeconds = 1 * ExpirationTime.Day)]
public async Task<List<SteamGameResponse>> GetGamesAsync(...)
  • 查询结果缓存 1 天

缓存失效

[InvalidateCache(Methods = [nameof(GetGamesAsync), nameof(GetGameAsync)])]
public async Task UpdateGamesAsync(...)
  • 更新数据时自动清除相关缓存
  • 通过 AOP 实现(需要配合拦截器框架)

7. 增量更新机制

private async Task UpdateGameAsync(List<SteamGame> dbGames, SteamGame updateGame)
{
    var dic = dbGame.UpdateContent(updateGame);
    if (dic.Count > 0) {
        logger.LogInformation("更新细节:{@UpdateContent}", dic);
    }
}
  • 只更新变化的字段
  • 记录详细的变更内容用于审计

设计亮点

  1. 防御性编程:配置校验、空值检查、异常重试
  2. 分布式友好:实例锁、幂等性设计
  3. 性能优化:并行处理、延迟加载、智能缓存
  4. 可观测性:详细的日志记录(包含实例 ID、游戏 ID 等上下文)
  5. 扩展性:事件驱动的资源处理机制

潜在改进点

  1. Parallel.ForEachAsync 可配置并发度(当前硬编码)
  2. API 延迟时间可配置化
  3. 批量更新的注释代码可能需要恢复(大量新游戏场景)
  4. 错误重试策略可进一步细化(如指数退避)
评论加载中...