ASP.NET核心 -- -- 第二部分:实用实例 (中文 (Chinese Simplified))

ASP.NET核心 -- -- 第二部分:实用实例

Thursday, 27 November 2025

//

12 minute read

理论是一回事; 生产代码是另一回事。 在第一部分中, 我们覆盖了抽象概念—— 现在让我们看看它们是如何应用到一个真正的代码库中。 这篇文章从这个博客平台上浏览了真实的背景服务: Polly 重新研究政策的文件观察者, 使用断路器的基于频道的电子邮件队列, 使用散列变化检测的语义搜索索引等等 。

一. 导言 导言 导言 导言 导言 导言 一,导言 导言 导言 导言 导言 导言

第一部分 第一部分我们探讨了实施ASP.NET核心背景服务的基本方法。 IHostedService, BackgroundService寿命周期管理, 以及关闭处理的常见陷阱。

现在是时候看到这些模式在起作用了。在这个文章中,我们将从一个制作博客平台上审视真实的背景服务, 展示:

  • 将标记下载文件同步到数据库的系统监视器
  • Polly 重试政策基于频道的电子邮件队列
  • 批量背景要求的分析事件发送器
  • 带有基于散列的变化探测的语义搜索索引器
  • 定期验证外部 URL 外部 URL 的断断链接检查器
  • 受扶养事务部门间的启动协调
  • 范围界定和实体实体框架生命周期管理

每个例子都说明了在建立生产背景服务时你将遇到的共同问题的实际解决办法。

示例1: 带有 Polly 重试的文件系统监视器

我们将检查的第一个服务, 将查看一个标记博客日志目录, 并在它们改变时自动处理它们。 这是由事件驱动的背景处理的完美例子 。

问题

创建或修改标记文件时 :

  1. 读取文件内容
  2. 分析分计元数据(标题、类别、公布日期)
  3. 输入 HTML 并提取纯文本
  4. 保存到数据库
  5. 触发翻译成其他语言
  6. 语义搜索索引

挑战: FileSystemWatcher 在文件仍在写入过程中, 事件可以点火, 导致 IOException 当你试着读的时候

解决方案: Markdown 董事会观察服务

public class MarkdownDirectoryWatcherService(
    MarkdownConfig markdownConfig,
    IServiceScopeFactory serviceScopeFactory,
    IStartupCoordinator startupCoordinator,
    ILogger<MarkdownDirectoryWatcherService> logger)
    : IHostedService
{
    private FileSystemWatcher _fileSystemWatcher;
    private Task _awaitChangeTask = Task.CompletedTask;

    public Task StartAsync(CancellationToken cancellationToken)
    {
        _fileSystemWatcher = new FileSystemWatcher
        {
            Path = markdownConfig.MarkdownPath,
            NotifyFilter = NotifyFilters.FileName | NotifyFilters.LastWrite |
                          NotifyFilters.CreationTime | NotifyFilters.Size,
            Filter = "*.md",
            IncludeSubdirectories = true
        };

        _fileSystemWatcher.EnableRaisingEvents = true;

        // Start background processing in a separate task
        _awaitChangeTask = Task.Run(() => AwaitChanges(cancellationToken), cancellationToken);

        logger.LogInformation("Started watching directory {Directory}", markdownConfig.MarkdownPath);

        // Signal ready - watcher is set up and listening
        startupCoordinator.SignalReady(StartupServiceNames.MarkdownDirectoryWatcher);

        return Task.CompletedTask;
    }

    public Task StopAsync(CancellationToken cancellationToken)
    {
        // Proper cleanup
        _fileSystemWatcher.EnableRaisingEvents = false;
        _fileSystemWatcher.Dispose();

        logger.LogInformation("Stopped watching directory");

        return Task.CompletedTask;
    }

    private async Task AwaitChanges(CancellationToken cancellationToken)
    {
        while (!cancellationToken.IsCancellationRequested)
        {
            var fileEvent = _fileSystemWatcher.WaitForChanged(WatcherChangeTypes.All);

            if (fileEvent.ChangeType == WatcherChangeTypes.Changed ||
                fileEvent.ChangeType == WatcherChangeTypes.Created)
            {
                await OnChangedAsync(fileEvent);
            }
            else if (fileEvent.ChangeType == WatcherChangeTypes.Deleted)
            {
                await OnDeletedAsync(fileEvent);
            }
            else if (fileEvent.ChangeType == WatcherChangeTypes.Renamed)
            {
                await OnRenamedAsync(fileEvent);
            }
        }
    }
}

注意几个关键模式 :

  1. Islodd Service, 而非背景服务 - 我们需要精细控制监视器的设置
  2. 单独背景任务 - 那个 StartAsync 立即返回返回,处理发生于 AwaitChanges
  3. 启动协调 - 准备接受依赖服务时发出信号
  4. 适当处置 - 禁用和处置文件系统监视器 StopAsync

处理 Polly 文件锁定问题

最有趣的部分是我们如何处理仍在写中的文件。 波利 是一个.NET抗御力图书馆, 提供回溯政策、断路器, 以及更多 :

private async Task OnChangedAsync(WaitForChangedResult e)
{
    if (e.Name == null) return;

    // Serilog activity for distributed tracing
    using var activity = Log.Logger.StartActivity("Markdown File Changed {Name}", e.Name);

    // Define a retry policy for file access issues
    var retryPolicy = Policy
        .Handle<IOException>() // Only handle IO exceptions (like file in use)
        .WaitAndRetryAsync(5, retryAttempt => TimeSpan.FromMilliseconds(500 * retryAttempt),
            (exception, timeSpan, retryCount, context) =>
            {
                activity?.Activity?.SetTag("Retry Attempt", retryCount);
                logger.LogWarning(
                    "File is in use, retrying attempt {RetryCount} after {TimeSpan}",
                    retryCount, timeSpan);
            });

    try
    {
        var fileName = e.Name;
        var isTranslated = Path.GetFileNameWithoutExtension(e.Name).Contains(".");
        var language = MarkdownBaseService.EnglishLanguage;
        var directory = markdownConfig.MarkdownPath;

        if (isTranslated)
        {
            language = Path.GetFileNameWithoutExtension(e.Name).Split('.').Last();
            fileName = Path.GetFileName(fileName);
            directory = markdownConfig.MarkdownTranslatedPath;
        }

        var filePath = Path.Combine(directory, fileName);

        using var scope = serviceScopeFactory.CreateScope();
        var blogService = scope.ServiceProvider.GetRequiredService<IBlogService>();

        // Use the Polly retry policy
        await retryPolicy.ExecuteAsync(async () =>
        {
            // Read the file - might throw IOException if locked
            var markdown = await File.ReadAllTextAsync(filePath);

            var slug = Path.GetFileNameWithoutExtension(fileName);
            if (isTranslated)
            {
                slug = slug.Split('.').First();
            }

            // Save to database
            var savedModel = await blogService.SavePost(slug, language, markdown);

            activity?.Activity?.SetTag("Page Processed", savedModel.Slug);

            // Index in semantic search (only for main directory files)
            if (!e.Name.Contains(Path.DirectorySeparatorChar))
            {
                await IndexPostForSemanticSearchAsync(scope, savedModel, language);
            }

            // Trigger translation for English posts
            if (language == MarkdownBaseService.EnglishLanguage &&
                !string.IsNullOrEmpty(savedModel.Markdown))
            {
                var translateService = scope.ServiceProvider
                    .GetRequiredService<IBackgroundTranslateService>();
                await translateService.TranslateForAllLanguages(
                    new PageTranslationModel
                    {
                        OriginalFileName = filePath,
                        OriginalMarkdown = savedModel.Markdown,
                        Persist = true
                    });
            }
        });

        activity?.Complete();
    }
    catch (Exception exception)
    {
        activity?.Complete(LogEventLevel.Error, exception);
    }
}

这里的密钥模式 :

  1. 投票重试政策 - 指数后退(500米、1秒、1.5秒、2秒、2.5秒)
  2. 仅重试 IO 例外 - 其他例外 泡沫起来
  3. 范围服务 - 建立每个文件的覆盖范围,以避免存在周期问题
  4. 结构性伐木 - 使用 血清 追踪和追踪活动的活动
  5. 连连业务 - 保存 索引 翻译

重试策略处理常有的情况, 当更改事件火灾时文本编辑器仍在写文件 。

视觉流动

graph TD
    A[File Modified] --> B[FileSystemWatcher Event]
    B --> C{File Locked?}
    C -->|Yes| D[Wait 500ms * Attempt]
    D --> E{Retry < 5?}
    E -->|Yes| C
    E -->|No| F[Log Error]
    C -->|No| G[Read File]
    G --> H[Parse Markdown]
    H --> I[Save to Database]
    I --> J{Main Directory?}
    J -->|Yes| K[Index for Search]
    J -->|No| L{English?}
    K --> L
    L -->|Yes| M[Trigger Translation]
    L -->|No| N[Complete]
    M --> N

    style A stroke:#059669,stroke-width:3px,color:#10b981
    style I stroke:#2563eb,stroke-width:3px,color:#3b82f6
    style K stroke:#7c3aed,stroke-width:3px,color:#8b5cf6
    style M stroke:#d97706,stroke-width:3px,color:#f59e0b

登记在 Program.cs (所有托管服务都采用这种模式):

builder.Services.AddHostedService<MarkdownDirectoryWatcherService>();

见《大会正式记录,2003年, Mostlylucid/Blog/WatcherService/MarkdownDirectoryWatcherService.cs.

实例2:基于频道的电子邮件队列

电子邮件发送是背景服务的一个典型使用实例。 您在谈判 SMTP 时不想阻拦 HTTP 请求, 所以您会排队在电子邮件上, 并在背景中发送 。

问题

当用户提交评论或联系表时:

  1. 审定和保存提交书
  2. 设置电子邮件通知
  3. 立即回复( 不要阻止 SMTP)
  4. 以重试逻辑在背景中发送邮件
  5. 优优优处理 SMTP 失败( 断路器)

解决方案:电子邮件SenderHostedService

此服务使用 a Channel<T> 用于排队和 Polly 抗御能力:

public class EmailSenderHostedService : IEmailSenderHostedService
{
    private readonly Channel<BaseEmailModel> _mailMessages =
        Channel.CreateUnbounded<BaseEmailModel>();
    private readonly CancellationTokenSource _cancellationTokenSource = new();
    private Task _sendTask = Task.CompletedTask;
    private readonly IEmailService _emailService;
    private readonly ILogger<EmailSenderHostedService> _logger;
    private readonly IAsyncPolicy _policyWrap;

    public EmailSenderHostedService(
        IEmailService emailService,
        ILogger<EmailSenderHostedService> logger)
    {
        _emailService = emailService;
        _logger = logger;

        // Retry policy: 3 attempts with exponential backoff
        var retryPolicy = Policy
            .Handle<SmtpException>()
            .WaitAndRetryAsync(3,
                attempt => TimeSpan.FromSeconds(2 * attempt),
                (exception, timeSpan, retryCount, context) =>
                {
                    logger.LogWarning(exception,
                        "Retry {RetryCount} for sending email failed", retryCount);
                });

        // Circuit breaker: open after 5 failures, stay open for 1 minute
        var circuitBreakerPolicy = Policy
            .Handle<SmtpException>()
            .CircuitBreakerAsync(
                5,
                TimeSpan.FromMinutes(1),
                onBreak: (exception, timespan) =>
                {
                    logger.LogError(
                        "Circuit broken due to too many failures. Breaking for {BreakDuration}",
                        timespan);
                },
                onReset: () =>
                {
                    logger.LogInformation("Circuit reset. Resuming email delivery.");
                },
                onHalfOpen: () =>
                {
                    logger.LogInformation("Circuit in half-open state. Testing connection...");
                });

        // Combine retry and circuit breaker
        _policyWrap = Policy.WrapAsync(retryPolicy, circuitBreakerPolicy);
    }

    public async Task SendEmailAsync(BaseEmailModel message)
    {
        await _mailMessages.Writer.WriteAsync(message);
    }

    public Task StartAsync(CancellationToken cancellationToken)
    {
        _logger.LogInformation("Starting background e-mail delivery");
        _sendTask = DeliverAsync(_cancellationTokenSource.Token);
        return Task.CompletedTask;
    }

    public async Task StopAsync(CancellationToken cancellationToken)
    {
        _logger.LogInformation("Stopping background e-mail delivery");

        // Proper shutdown sequence
        await _cancellationTokenSource.CancelAsync();
        _mailMessages.Writer.Complete(); // Critical: complete the channel

        // Wait for the background task to finish
        await Task.WhenAny(_sendTask, Task.Delay(Timeout.Infinite, cancellationToken));
    }

    private async Task DeliverAsync(CancellationToken token)
    {
        _logger.LogInformation("E-mail background delivery started");

        try
        {
            // Process items as they arrive
            while (await _mailMessages.Reader.WaitToReadAsync(token))
            {
                BaseEmailModel? message = null;
                try
                {
                    message = await _mailMessages.Reader.ReadAsync(token);

                    // Execute with retry policy and circuit breaker
                    await _policyWrap.ExecuteAsync(async () =>
                    {
                        switch (message)
                        {
                            case ContactEmailModel contactEmailModel:
                                await _emailService.SendContactEmail(contactEmailModel);
                                break;
                            case CommentEmailModel commentEmailModel:
                                await _emailService.SendCommentEmail(commentEmailModel);
                                break;
                            case ConfirmEmailModel confirmEmailModel:
                                await _emailService.SendConfirmationEmail(confirmEmailModel);
                                break;
                        }
                    });

                    _logger.LogInformation("Email from {SenderEmail} sent", message.SenderEmail);
                }
                catch (OperationCanceledException)
                {
                    break; // Shutdown requested
                }
                catch (Exception exc)
                {
                    _logger.LogError(exc,
                        "Couldn't send an e-mail from {SenderEmail}",
                        message?.SenderEmail);
                }
            }
        }
        catch (OperationCanceledException)
        {
            _logger.LogWarning("E-mail background delivery canceled");
        }

        _logger.LogInformation("E-mail background delivery stopped");
    }

    public void Dispose()
    {
        _cancellationTokenSource.Dispose();
    }
}

所展示的关键模式:

  1. 频道作为队列 - 无约束的电文频道
  2. 政策总结 - 在断路器内重试
  3. 优优优关闭 - 完成频道、取消标记、等待任务
  4. 模式匹配模式匹配 - 开启不同的电子邮件类型
  5. 涵盖寿命期 - 曾作为单子建立过一次, 保持长寿状态

重试和电路断路器可视化

graph TD
    A[Queue Email] --> B[Write to Channel]
    B --> C[Background Task Reads]
    C --> D{Send Email}
    D -->|Success| E[Log Success]
    D -->|SmtpException| F{Retry Count < 3?}
    F -->|Yes| G[Wait 2s * Attempt]
    G --> D
    F -->|No| H{Circuit Breaker}
    H -->|< 5 Failures| I[Log Failure]
    H -->|5+ Failures| J[Open Circuit]
    J --> K[Wait 1 Minute]
    K --> L[Half-Open: Test]
    L -->|Success| M[Close Circuit]
    L -->|Failure| J
    E --> N[Continue]
    I --> N

    style A stroke:#059669,stroke-width:3px,color:#10b981
    style D stroke:#2563eb,stroke-width:3px,color:#3b82f6
    style J stroke:#dc2626,stroke-width:3px,color:#ef4444
    style M stroke:#059669,stroke-width:3px,color:#10b981

这表明生产等级模式:

  • 重试政策 处理瞬态故障(临时网络问题)
  • 断路器 (如果SMTP被击落,不要一直敲门)
  • 频道频道 提供自然后压( 如果我们无法发送, 信件排队)

实例3:分析事件发件人(乌干达)

此服务性队队列分析事件, 并发送到 木美 分析服务器。 它类似于电子邮件服务, 但更简单,不需要重试政策, 只需要射击和忘记。

public class UmamiBackgroundSender(
    IServiceScopeFactory scopeFactory,
    ILogger<UmamiBackgroundSender> logger) : IHostedService
{
    private readonly CancellationTokenSource _cancellationTokenSource = new();
    private readonly Channel<SendBackgroundPayload> _channel =
        Channel.CreateUnbounded<SendBackgroundPayload>();
    private Task _sendTask = Task.CompletedTask;

    public Task StartAsync(CancellationToken cancellationToken)
    {
        _sendTask = SendRequest(_cancellationTokenSource.Token);
        return Task.CompletedTask;
    }

    public async Task StopAsync(CancellationToken cancellationToken)
    {
        logger.LogInformation("UmamiBackgroundSender is stopping.");

        // Standard shutdown pattern
        await _cancellationTokenSource.CancelAsync();
        _channel.Writer.Complete();

        try
        {
            await Task.WhenAny(_sendTask, Task.Delay(Timeout.Infinite, cancellationToken));
        }
        catch (OperationCanceledException)
        {
            logger.LogWarning("StopAsync operation was canceled.");
        }
    }

    public async Task TrackPageView(string url, string title, UmamiPayload? payload = null)
    {
        await using var scope = scopeFactory.CreateAsyncScope();
        var payloadService = scope.ServiceProvider.GetRequiredService<PayloadService>();
        var sendPayload = payloadService.PopulateFromPayload(payload, null);
        sendPayload.Url = url;
        sendPayload.Title = title;

        await _channel.Writer.WriteAsync(new SendBackgroundPayload("event", sendPayload));
        logger.LogInformation("Umami pageview event queued");
    }

    private async Task SendRequest(CancellationToken token)
    {
        logger.LogInformation("Umami background delivery started");

        // Double while loop: outer waits for items, inner drains all available
        while (await _channel.Reader.WaitToReadAsync(token))
        {
            while (_channel.Reader.TryRead(out var payload))
            {
                try
                {
                    using var scope = scopeFactory.CreateScope();
                    var client = scope.ServiceProvider.GetRequiredService<UmamiClient>();

                    await client.Send(payload.Payload, type: payload.EventType);

                    logger.LogInformation("Umami background event sent: {EventType}",
                        payload.EventType);
                }
                catch (OperationCanceledException)
                {
                    logger.LogWarning("Umami background delivery canceled.");
                    return;
                }
                catch (Exception ex)
                {
                    logger.LogError(ex, "Error sending Umami background event.");
                }
            }
        }
    }

    private record SendBackgroundPayload(string EventType, UmamiPayload Payload);
}

电子邮件服务的关键差异 :

  1. 无复原力政策 - 分析不是关键,如果它失败了,我们只要记录和继续
  2. 双环双循环 - 内部循环排水所有可用项目,以提高效率
  3. 范围范围服务创建 - 每个发送器都为 HTTP 客户提供新的视野

这表明并非所有背景服务都需要复杂的错误处理。 对于非临界遥测而言,简单的记录可能就足够了。

实例4:与背景服务有关的定期背景工作

现在让我们来看看那些定期工作而不是处理队列的服务。 BrokenLinkCheckerBackgroundService 校对:Portnoy

public class BrokenLinkCheckerBackgroundService : BackgroundService
{
    private readonly IServiceProvider _serviceProvider;
    private readonly ILogger<BrokenLinkCheckerBackgroundService> _logger;
    private readonly HttpClient _httpClient;
    private readonly TimeSpan _checkInterval = TimeSpan.FromHours(1);
    private const int BatchSize = 20;

    public BrokenLinkCheckerBackgroundService(
        IServiceProvider serviceProvider,
        ILogger<BrokenLinkCheckerBackgroundService> logger,
        IHttpClientFactory httpClientFactory)
    {
        _serviceProvider = serviceProvider;
        _logger = logger;
        _httpClient = httpClientFactory.CreateClient("BrokenLinkChecker");
        _httpClient.Timeout = TimeSpan.FromSeconds(30);
    }

    protected override async Task ExecuteAsync(CancellationToken stoppingToken)
    {
        _logger.LogInformation("Broken Link Checker Background Service started");

        // Initial delay to let the application fully start
        await Task.Delay(TimeSpan.FromMinutes(5), stoppingToken);

        while (!stoppingToken.IsCancellationRequested)
        {
            try
            {
                await CheckLinksAsync(stoppingToken);
                await FetchArchiveUrlsAsync(stoppingToken);
            }
            catch (Exception ex)
            {
                _logger.LogError(ex, "Error in Broken Link Checker Background Service");
            }

            _logger.LogInformation("Broken Link Checker sleeping for {Interval}", _checkInterval);
            await Task.Delay(_checkInterval, stoppingToken);
        }

        _logger.LogInformation("Broken Link Checker Background Service stopped");
    }

    private async Task CheckLinksAsync(CancellationToken cancellationToken)
    {
        _logger.LogInformation("Starting link validity check");

        using var scope = _serviceProvider.CreateScope();
        var brokenLinkService = scope.ServiceProvider.GetRequiredService<IBrokenLinkService>();

        var linksToCheck = await brokenLinkService.GetLinksToCheckAsync(BatchSize, cancellationToken);
        _logger.LogInformation("Found {Count} links to check", linksToCheck.Count);

        foreach (var link in linksToCheck)
        {
            if (cancellationToken.IsCancellationRequested) break;

            try
            {
                var (statusCode, isBroken, error) = await CheckUrlAsync(link.OriginalUrl, cancellationToken);
                await brokenLinkService.UpdateLinkStatusAsync(
                    link.Id, statusCode, isBroken, error, cancellationToken);

                if (isBroken)
                {
                    _logger.LogWarning("Link is broken: {Url} (Status: {StatusCode})",
                        link.OriginalUrl, statusCode);
                }

                // Be respectful to servers
                await Task.Delay(TimeSpan.FromSeconds(2), cancellationToken);
            }
            catch (Exception ex)
            {
                _logger.LogError(ex, "Error checking link: {Url}", link.OriginalUrl);
                await brokenLinkService.UpdateLinkStatusAsync(
                    link.Id, 0, true, ex.Message, cancellationToken);
            }
        }
    }

    private async Task<(int statusCode, bool isBroken, string? error)> CheckUrlAsync(
        string url,
        CancellationToken cancellationToken)
    {
        try
        {
            using var request = new HttpRequestMessage(HttpMethod.Head, url);
            request.Headers.UserAgent.ParseAdd(
                "Mozilla/5.0 (compatible; MostlylucidBot/1.0; +https://www.mostlylucid.net)");

            using var response = await _httpClient.SendAsync(
                request,
                HttpCompletionOption.ResponseHeadersRead,
                cancellationToken);

            var statusCode = (int)response.StatusCode;
            var isBroken = response.StatusCode == HttpStatusCode.NotFound ||
                          response.StatusCode == HttpStatusCode.Gone ||
                          statusCode >= 500;

            return (statusCode, isBroken, null);
        }
        catch (HttpRequestException ex)
        {
            return (0, true, ex.Message);
        }
        catch (TaskCanceledException ex) when (ex.InnerException is TimeoutException)
        {
            return (0, true, "Request timed out");
        }
    }
}

显示的模式显示:

  1. 背景服务基级 - 适合定期工作
  2. 初始延迟 - 等待其它服务到初始端
  3. 条纹 - 一次处理20个链接,不是一次全部处理
  4. 限 限 限 利率 - 为礼貌起见,两次检查之间延迟2秒
  5. 范围范围服务创建 - 每批新范围
  6. 取消代号检查 - 如果要求停工,请提前停工

实例5:用基于散装物的变更检测法的语义搜索索引

语义搜索索引器更复杂 — — 只能是实际改变的重新索引。 这节省了昂贵的嵌入 API 电话 。

csharp public class SemanticIndexingBackgroundService : BackgroundService { private readonly ILogger _logger; private readonly IServiceProvider _serviceProvider; private readonly ISemanticSearchService _semanticSearchService; private readonly MarkdownConfig _markdownConfig; private readonly SemanticSearchConfig _semanticSearchConfig; private readonly TimeSpan _indexInterval = TimeSpan.FromHours(1); private readonly TimeSpan _startupDelay = TimeSpan.FromSeconds(30);

protected override async Task ExecuteAsync(CancellationToken stoppingToken)
{
    if (!_semanticSearchConfig.Enabled)
    {
        _logger.LogInformation("Semantic search is disabled, indexing service will not run");
        return;
    }

    _logger.LogInformation("Semantic indexing background service starting...");

    // Wait for other services to initialise
    await Task.Delay(_startupDelay, stoppingToken);

    // Initialise the semantic search service
    try
    {
        await _semanticSearchService.InitializeAsync(stoppingToken);
        _logger.LogInformation("Semantic search initialized successfully");
    }
    catch (Exception ex)
    {
        _logger.LogError(ex, "Failed to initialize semantic search service");
        return;
    }

    // Initial indexing
    await IndexAllMarkdownFilesAsync(stoppingToken);

    // Periodic re-indexing to catch any changes
    while (!stoppingToken.IsCancellationRequested)
    {
        try
        {
            await Task.Delay(_indexInterval, stoppingToken);
            await IndexAllMarkdownFilesAsync(stoppingToken);
        }
        catch (OperationCanceledException)
        {
            break;
        }
        catch (Exception ex)
        {
            _logger.LogError(ex, "Error during periodic indexing");
        }
    }

    _logger.LogInformation("Semantic indexing background service stopped");
}

private async Task IndexAllMarkdownFilesAsync(CancellationToken stoppingToken)
{
    var markdownPath = _markdownConfig.MarkdownPath;

    if (!Directory.Exists(markdownPath))
    {
        _logger.LogWarning("Markdown directory does not exist: {Path}", markdownPath);
        return;
    }

    // Get ONLY files in the main directory, NOT subdirectories
    var markdownFiles = Directory.GetFiles(markdownPath, "*.md", SearchOption.TopDirectoryOnly);

    _logger.LogInformation("Found {Count} markdown files to index in {Path}",
        markdownFiles.Length, markdownPath);

    var indexedCount = 0;
    var skippedCount = 0;
    var errorCount = 0;

    using var scope = _serviceProvider.CreateScope();
    var markdownRenderingService = scope.ServiceProvider
        .GetRequiredService<MarkdownRenderingService>();

    foreach (var filePath in markdownFiles)
    {
        if (stoppingToken.IsCancellationRequested)
            break;

        try
        {
            var result = await IndexMarkdownFileAsync(
                filePath,
                markdownRenderingService,
                stoppingToken);

            if (result == IndexResult.Indexed)
                indexedCount++;
            else if (result == IndexResult.Skipped)
                skippedCount++;
        }
        catch (Exception ex)
        {
            errorCount++;
            _logger.LogError(ex, "Error indexing file: {FilePath}", filePath);
        }

        // Delay to avoid overwhelming the embedding service
        await Task.Delay(100, stoppingToken);
    }

    _logger.LogInformation(
        "Indexing complete: {Indexed} indexed, {Skipped} skipped (unchanged), {Errors} errors",
        indexedCount, skippedCount, errorCount);
}

private async Task<IndexResult> IndexMarkdownFileAsync(
    string filePath,
    MarkdownRenderingService markdownRenderingService,
    CancellationToken stoppingToken)
{
    var fileName = Path.GetFileNameWithoutExtension(filePath);

    // Skip translated files (they have language suffix like .es.md, .fr.md)
    if (fileName.Contains('.'))
    {
        var parts = fileName.Split('.');
        if (parts.Length >= 2 && parts[^1].Length == 2)
        {
            return IndexResult.Skipped;
        }
    }

    var markdown = await File.ReadAllTextAsync(filePath, stoppingToken);
    var fileInfo = new FileInfo(filePath);

    var blogPost = markdownRenderingService.GetPageFromMarkdown(
        markdown,
        fileInfo.LastWriteTimeUtc,
        filePath);

    if (blogPost.IsHidden)
    {
        _logger.LogDebug("Skipping hidden post: {Slug}", blogPost.Slug);
        return IndexResult.Skipped;
    }

    // Compute content hash
    var contentHash = ComputeContentHash(blogPost.PlainTextContent);

    // Check if reindexing is needed
    var needsReindex = await _semanticSearchService.NeedsReindexingAsync(
        blogPost.Slug,
        MarkdownBaseService.EnglishLanguage,
        contentHash,
        stoppingToken);

    if (!needsReindex)
    {
        _logger.LogDebug("Skipping unchanged post: {Slug}", blogPost.Slug);
        return IndexResult.Skipped;
    }

    // Create document for indexing
    var document = new BlogPostDocument
    {
        Id = $"{blogPost.Slug}_{MarkdownBaseService.EnglishLanguage}",
        Slug = blogPost.Slug,
        Title = blogPost.Title,
        Content = blogPost.Plain
Finding related posts...
logo

© 2026 Scott Galloway — Unlicense — All content and source code on this site is free to use, copy, modify, and sell.