Skip to content

Commit ff8ecd1

Browse files
feat: implement async background post content extraction to save HTML in blob storage (#185)
Co-authored-by: fboucher-os <fboucheros+github@gmail.com>
1 parent 5077f44 commit ff8ecd1

9 files changed

Lines changed: 290 additions & 1 deletion

File tree

Lines changed: 74 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,74 @@
1+
using FluentAssertions;
2+
using Microsoft.Extensions.DependencyInjection;
3+
using NoteBookmark.Api.Tests.Fixtures;
4+
using NoteBookmark.Domain;
5+
using System;
6+
using System.Net;
7+
using System.Net.Http;
8+
using System.Net.Http.Json;
9+
using System.Threading.Tasks;
10+
using Xunit;
11+
using Azure.Storage.Blobs;
12+
13+
namespace NoteBookmark.Api.Tests.Endpoints;
14+
15+
public class PostExtractionTests : IClassFixture<NoteBookmarkApiTestFactory>
16+
{
17+
private readonly NoteBookmarkApiTestFactory _factory;
18+
private readonly HttpClient _client;
19+
20+
public PostExtractionTests(NoteBookmarkApiTestFactory factory)
21+
{
22+
_factory = factory;
23+
_client = _factory.CreateClient();
24+
}
25+
26+
[Fact]
27+
public async Task ExtractPostDetails_TriggersBackgroundWorkerAndSavesHtmlToBlobStorage()
28+
{
29+
// Arrange
30+
var url = "https://example.com/blog/test-post-" + Guid.NewGuid();
31+
var extractRequest = new
32+
{
33+
url = url,
34+
tags = "test",
35+
category = "Test"
36+
};
37+
38+
// Act - Call the API to extract metadata and save the post
39+
var response = await _client.PostAsJsonAsync("/api/posts/extractPostDetails", extractRequest);
40+
41+
// Assert API response is OK
42+
response.StatusCode.Should().Be(HttpStatusCode.OK);
43+
44+
var post = await response.Content.ReadFromJsonAsync<Post>();
45+
post.Should().NotBeNull();
46+
var postId = post!.Id ?? post.RowKey;
47+
postId.Should().NotBeNullOrEmpty();
48+
49+
// Since the extraction happens asynchronously in a BackgroundWorker,
50+
// we poll the blob storage for a short time to verify the file was created.
51+
var blobServiceClient = _factory.Services.GetRequiredService<BlobServiceClient>();
52+
var containerClient = blobServiceClient.GetBlobContainerClient("cleanedposts");
53+
var blobClient = containerClient.GetBlobClient($"{postId}.html");
54+
55+
// Wait up to 5 seconds for the background worker to process
56+
bool blobExists = false;
57+
for (int i = 0; i < 25; i++)
58+
{
59+
if (await blobClient.ExistsAsync())
60+
{
61+
blobExists = true;
62+
break;
63+
}
64+
await Task.Delay(200);
65+
}
66+
67+
blobExists.Should().BeTrue("HTML content should be processed by the background worker and saved to blob storage");
68+
69+
// Verify the content saved matches the fake content
70+
var downloadResult = await blobClient.DownloadContentAsync();
71+
var content = downloadResult.Value.Content.ToString();
72+
content.Should().Contain(url);
73+
}
74+
}
Lines changed: 13 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,13 @@
1+
using System.Threading;
2+
using System.Threading.Tasks;
3+
4+
namespace NoteBookmark.Api.Tests.Fixtures;
5+
6+
public class FakePostParserClient : IPostParserClient
7+
{
8+
public Task<string?> ExtractContentAsync(string url, CancellationToken cancellationToken = default)
9+
{
10+
// Return a mock HTML snippet for testing
11+
return Task.FromResult<string?>($"<div>Extracted HTML content for {url}</div>");
12+
}
13+
}

‎src/NoteBookmark.Api.Tests/Fixtures/NoteBookmarkApiTestFactory.cs‎

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -34,6 +34,9 @@ protected override void ConfigureWebHost(IWebHostBuilder builder)
3434
services.AddSingleton(new TableServiceClient(connectionString));
3535
services.AddSingleton(new BlobServiceClient(connectionString));
3636
}
37+
38+
// Register FakePostParserClient for integration tests
39+
services.AddSingleton<NoteBookmark.Api.IPostParserClient, FakePostParserClient>();
3740
});
3841
}
3942

Lines changed: 9 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,9 @@
1+
using System.Threading;
2+
using System.Threading.Tasks;
3+
4+
namespace NoteBookmark.Api;
5+
6+
public interface IPostParserClient
7+
{
8+
Task<string?> ExtractContentAsync(string url, CancellationToken cancellationToken = default);
9+
}

‎src/NoteBookmark.Api/PostEndpoints.cs‎

Lines changed: 9 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -94,7 +94,11 @@ static Results<Ok, BadRequest> SavePost(Post post, TableServiceClient tblClient,
9494
}
9595
return TypedResults.BadRequest();
9696
}
97-
static async Task<Results<Ok<Post>, BadRequest>> ExtractPostDetails(ExtractPostRequest request, TableServiceClient tblClient, BlobServiceClient blobClient)
97+
static async Task<Results<Ok<Post>, BadRequest>> ExtractPostDetails(
98+
ExtractPostRequest request,
99+
TableServiceClient tblClient,
100+
BlobServiceClient blobClient,
101+
PostExtractionQueue queue)
98102
{
99103
var dataStorageService = new DataStorageService(tblClient, blobClient);
100104

@@ -105,6 +109,10 @@ static async Task<Results<Ok<Post>, BadRequest>> ExtractPostDetails(ExtractPostR
105109
if (post != null)
106110
{
107111
dataStorageService.SavePost(post);
112+
113+
// Queue background HTML extraction task
114+
queue.QueueBackgroundWorkItem(new ExtractionTask(post.Id ?? post.RowKey, post.Url ?? decodeUrl));
115+
108116
return TypedResults.Ok(post);
109117
}
110118
return TypedResults.BadRequest();
Lines changed: 87 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,87 @@
1+
using System;
2+
using System.IO;
3+
using System.Text;
4+
using System.Threading;
5+
using System.Threading.Tasks;
6+
using Azure.Storage.Blobs;
7+
using Microsoft.Extensions.DependencyInjection;
8+
using Microsoft.Extensions.Hosting;
9+
using Microsoft.Extensions.Logging;
10+
11+
namespace NoteBookmark.Api;
12+
13+
public class PostExtractionBackgroundWorker : BackgroundService
14+
{
15+
private readonly PostExtractionQueue _queue;
16+
private readonly IServiceProvider _serviceProvider;
17+
private readonly ILogger<PostExtractionBackgroundWorker> _logger;
18+
19+
public PostExtractionBackgroundWorker(
20+
PostExtractionQueue queue,
21+
IServiceProvider serviceProvider,
22+
ILogger<PostExtractionBackgroundWorker> logger)
23+
{
24+
_queue = queue;
25+
_serviceProvider = serviceProvider;
26+
_logger = logger;
27+
}
28+
29+
protected override async Task ExecuteAsync(CancellationToken stoppingToken)
30+
{
31+
_logger.LogInformation("Post Extraction Background Worker started.");
32+
33+
while (!stoppingToken.IsCancellationRequested)
34+
{
35+
try
36+
{
37+
var task = await _queue.DequeueAsync(stoppingToken);
38+
_logger.LogInformation("Processing extraction for Post: {PostId}, URL: {Url}", task.PostId, task.Url);
39+
40+
await ProcessExtractionAsync(task, stoppingToken);
41+
}
42+
catch (OperationCanceledException)
43+
{
44+
// Normal shutdown
45+
break;
46+
}
47+
catch (Exception ex)
48+
{
49+
_logger.LogError(ex, "Error occurred executing background extraction task.");
50+
}
51+
}
52+
53+
_logger.LogInformation("Post Extraction Background Worker stopped.");
54+
}
55+
56+
private async Task ProcessExtractionAsync(ExtractionTask task, CancellationToken cancellationToken)
57+
{
58+
using var scope = _serviceProvider.CreateScope();
59+
var parserClient = scope.ServiceProvider.GetRequiredService<IPostParserClient>();
60+
var blobServiceClient = scope.ServiceProvider.GetRequiredService<BlobServiceClient>();
61+
62+
try
63+
{
64+
var content = await parserClient.ExtractContentAsync(task.Url, cancellationToken);
65+
if (string.IsNullOrEmpty(content))
66+
{
67+
_logger.LogWarning("No content returned for URL: {Url}. Skipping blob upload.", task.Url);
68+
return;
69+
}
70+
71+
var containerClient = blobServiceClient.GetBlobContainerClient("cleanedposts");
72+
await containerClient.CreateIfNotExistsAsync(cancellationToken: cancellationToken);
73+
74+
var blobClient = containerClient.GetBlobClient($"{task.PostId}.html");
75+
76+
byte[] contentBytes = Encoding.UTF8.GetBytes(content);
77+
using var stream = new MemoryStream(contentBytes);
78+
79+
await blobClient.UploadAsync(stream, overwrite: true, cancellationToken: cancellationToken);
80+
_logger.LogInformation("Successfully saved extracted HTML for Post {PostId} to Blob Storage.", task.PostId);
81+
}
82+
catch (Exception ex)
83+
{
84+
_logger.LogError(ex, "Failed to process extraction for Post {PostId} / URL: {Url}", task.PostId, task.Url);
85+
}
86+
}
87+
}
Lines changed: 34 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,34 @@
1+
using System;
2+
using System.Threading;
3+
using System.Threading.Channels;
4+
using System.Threading.Tasks;
5+
6+
namespace NoteBookmark.Api;
7+
8+
public record ExtractionTask(string PostId, string Url);
9+
10+
public class PostExtractionQueue
11+
{
12+
private readonly Channel<ExtractionTask> _queue;
13+
14+
public PostExtractionQueue()
15+
{
16+
// Unbounded channel is simple and suitable for this task queue.
17+
_queue = Channel.CreateUnbounded<ExtractionTask>(new UnboundedChannelOptions
18+
{
19+
SingleReader = true,
20+
SingleWriter = false
21+
});
22+
}
23+
24+
public void QueueBackgroundWorkItem(ExtractionTask task)
25+
{
26+
ArgumentNullException.ThrowIfNull(task);
27+
_queue.Writer.TryWrite(task);
28+
}
29+
30+
public async ValueTask<ExtractionTask> DequeueAsync(CancellationToken cancellationToken)
31+
{
32+
return await _queue.Reader.ReadAsync(cancellationToken);
33+
}
34+
}
Lines changed: 56 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,56 @@
1+
using System;
2+
using System.Net.Http;
3+
using System.Net.Http.Json;
4+
using System.Text.Json.Serialization;
5+
using System.Threading;
6+
using System.Threading.Tasks;
7+
using Microsoft.Extensions.Logging;
8+
9+
namespace NoteBookmark.Api;
10+
11+
public class PostParserClient : IPostParserClient
12+
{
13+
private readonly HttpClient _httpClient;
14+
private readonly ILogger<PostParserClient> _logger;
15+
16+
public PostParserClient(HttpClient httpClient, ILogger<PostParserClient> logger)
17+
{
18+
_httpClient = httpClient;
19+
_logger = logger;
20+
// Configure base address or default headers if needed, but since URL is fully specified we can just configure it or call it directly.
21+
if (_httpClient.BaseAddress == null)
22+
{
23+
_httpClient.BaseAddress = new Uri("https://azpostlight-parser.azurewebsites.net/");
24+
}
25+
}
26+
27+
public async Task<string?> ExtractContentAsync(string url, CancellationToken cancellationToken = default)
28+
{
29+
try
30+
{
31+
_logger.LogInformation("Calling parser API for URL: {Url}", url);
32+
var requestBody = new { url = url };
33+
var response = await _httpClient.PostAsJsonAsync("parser", requestBody, cancellationToken);
34+
35+
if (!response.IsSuccessStatusCode)
36+
{
37+
_logger.LogWarning("Parser API returned error status: {StatusCode}", response.StatusCode);
38+
return null;
39+
}
40+
41+
var result = await response.Content.ReadFromJsonAsync<ParserResponse>(cancellationToken: cancellationToken);
42+
return result?.Content;
43+
}
44+
catch (Exception ex)
45+
{
46+
_logger.LogError(ex, "Failed to extract content for URL: {Url}", url);
47+
return null;
48+
}
49+
}
50+
51+
private class ParserResponse
52+
{
53+
[JsonPropertyName("content")]
54+
public string? Content { get; set; }
55+
}
56+
}

‎src/NoteBookmark.Api/Program.cs‎

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -15,6 +15,11 @@
1515
// Register data storage service
1616
builder.Services.AddScoped<IDataStorageService, DataStorageService>();
1717

18+
// Register background extraction queue and worker
19+
builder.Services.AddHttpClient<IPostParserClient, PostParserClient>();
20+
builder.Services.AddSingleton<PostExtractionQueue>();
21+
builder.Services.AddHostedService<PostExtractionBackgroundWorker>();
22+
1823
// Register AI settings provider
1924
builder.Services.AddScoped<IAISettingsProvider, AISettingsProvider>();
2025

0 commit comments

Comments
 (0)