From 3b93893618ab9d5cae65e51432b9b3502bd3d51d Mon Sep 17 00:00:00 2001 From: cesnimda Date: Thu, 2 Jul 2026 16:50:27 +0200 Subject: [PATCH] feat(ai): semantic-search backfill worker + Ollama compose profile (RECOMMENDATIONS #4a) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit First activation slice of semantic search (docs/discovery/05+06): - EmbeddingBackfillWorker: fills Email.Embedding in small batches (32, 2s pause, newest first; subject+snippet, truncated) so interactive load is never starved. Exits immediately when the embedding provider is unavailable — AI-off deployments are untouched, Ollama hiccups back off instead of crashing the host. - compose: optional 'ollama' service under the 'ai' profile (persistent model volume, commented GPU passthrough for the RTX 3080); api gets Ai__OllamaBaseUrl pointing at it. - LiveDb test proves the worker persists 768-dim vectors against real pgvector (deterministic fake provider); LiveDb classes serialised into one xUnit collection (concurrent MigrateAsync on a fresh DB races — found while testing). Verified against REAL Ollama locally: pulled nomic-embed-text in a container and confirmed the /api/embeddings contract the provider uses returns 768-dim vectors. Full suite 55 green (4 LiveDb vs real pgvector); format clean. Next slice: hybrid RRF ranking in SearchService once embeddings exist. Co-Authored-By: Claude Opus 4.8 --- docker-compose.yml | 20 ++++ .../Ai/EmbeddingBackfillWorker.cs | 104 ++++++++++++++++++ .../DependencyInjection.cs | 3 + .../InboxIntel.Infrastructure.csproj | 3 + .../EmbeddingBackfillTests.cs | 95 ++++++++++++++++ .../LiveDbSearchTests.cs | 1 + 6 files changed, 226 insertions(+) create mode 100644 src/InboxIntel.Infrastructure/Ai/EmbeddingBackfillWorker.cs create mode 100644 tests/InboxIntel.IntegrationTests/EmbeddingBackfillTests.cs diff --git a/docker-compose.yml b/docker-compose.yml index 6b370d9..59bc83c 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -34,6 +34,8 @@ services: GoogleOAuth__ClientId: ${GOOGLE_CLIENT_ID:-} GoogleOAuth__ClientSecret: ${GOOGLE_CLIENT_SECRET:-} Ai__Mode: ${AI_MODE:-Disabled} + # Points at the compose 'ollama' service when the ai profile is up; harmless otherwise. + Ai__OllamaBaseUrl: ${OLLAMA_BASE_URL:-http://ollama:11434} # Dev mode shows the dev banner and caps the initial sync. Set DEV_MODE=true # and MAX_MESSAGES=1000 in deploy/.env to exercise it in this Docker setup. App__DevMode: ${DEV_MODE:-false} @@ -59,6 +61,23 @@ services: ports: - "8081:80" + # Local AI (semantic search + assistants). Enable with: + # docker compose --profile ai up -d && set AI_MODE=LocalOllama in deploy/.env + # First run: docker compose exec ollama ollama pull nomic-embed-text + # GPU (RTX 3080): uncomment the deploy block to pass the GPU through. + ollama: + image: ollama/ollama + profiles: ["ai"] + volumes: + - ollama:/root/.ollama + # deploy: + # resources: + # reservations: + # devices: + # - driver: nvidia + # count: all + # capabilities: [gpu] + # Optional reverse proxy. Enable with: docker compose --profile proxy up nginx: image: nginx:alpine @@ -74,3 +93,4 @@ services: volumes: pgdata: keys: + ollama: diff --git a/src/InboxIntel.Infrastructure/Ai/EmbeddingBackfillWorker.cs b/src/InboxIntel.Infrastructure/Ai/EmbeddingBackfillWorker.cs new file mode 100644 index 0000000..682db3f --- /dev/null +++ b/src/InboxIntel.Infrastructure/Ai/EmbeddingBackfillWorker.cs @@ -0,0 +1,104 @@ +using InboxIntel.Application.Abstractions; +using InboxIntel.Infrastructure.Persistence; +using Microsoft.EntityFrameworkCore; +using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Hosting; +using Microsoft.Extensions.Logging; + +namespace InboxIntel.Infrastructure.Ai; + +/// +/// Fills Email.Embedding (pgvector) for semantic search, in small background batches +/// so interactive requests are never starved (per docs/discovery/06: embeddings are the +/// small always-on model; the batch pause keeps VRAM/CPU pressure low). Exits immediately +/// when the embedding provider is unavailable (AI disabled / Ollama down) — semantic search +/// simply stays dormant and lexical search is unaffected. +/// +public class EmbeddingBackfillWorker : BackgroundService +{ + private const int BatchSize = 32; + private static readonly TimeSpan BatchPause = TimeSpan.FromSeconds(2); + private static readonly TimeSpan IdleRescan = TimeSpan.FromMinutes(15); + + private readonly IServiceScopeFactory _scopeFactory; + private readonly ILogger _logger; + + public EmbeddingBackfillWorker(IServiceScopeFactory scopeFactory, ILogger logger) + { + _scopeFactory = scopeFactory; + _logger = logger; + } + + protected override async Task ExecuteAsync(CancellationToken stoppingToken) + { + // Provider availability is fixed by configuration for the process lifetime. + using (var probe = _scopeFactory.CreateScope()) + { + if (!probe.ServiceProvider.GetRequiredService().IsAvailable) + { + _logger.LogDebug("EmbeddingBackfillWorker idle: no embedding provider (AI disabled)."); + return; + } + } + + _logger.LogInformation("EmbeddingBackfillWorker started (batch {Batch}, pause {Pause}s)", + BatchSize, BatchPause.TotalSeconds); + + while (!stoppingToken.IsCancellationRequested) + { + int processed; + try + { + processed = await ProcessBatchAsync(stoppingToken); + } + catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested) { break; } + catch (Exception ex) + { + // Ollama hiccups must never crash the host; back off and retry. + _logger.LogWarning(ex, "Embedding batch failed; retrying after idle pause."); + processed = 0; + } + + await Task.Delay(processed > 0 ? BatchPause : IdleRescan, stoppingToken); + } + } + + /// Embeds one batch. Public-ish (internal) for direct testing. + internal async Task ProcessBatchAsync(CancellationToken ct) + { + using var scope = _scopeFactory.CreateScope(); + var db = scope.ServiceProvider.GetRequiredService(); + var embeddings = scope.ServiceProvider.GetRequiredService(); + + var batch = await db.Emails + .Where(e => e.Embedding == null) + .OrderByDescending(e => e.SentAtUtc) // newest mail becomes searchable first + .Take(BatchSize) + .ToListAsync(ct); + if (batch.Count == 0) return 0; + + // Subject + snippet is the semantic core; bodies are noisy (signatures, quoting) + // and slow to embed. Truncate defensively to keep well inside the model context. + var texts = batch + .Select(e => Truncate($"{e.Subject}\n{e.Snippet ?? e.BodyText}", 2000)) + .ToList(); + var vectors = await embeddings.EmbedBatchAsync(texts, ct); + if (vectors.Count != batch.Count) + { + _logger.LogWarning("Embedding batch returned {Got} vectors for {Want} emails; skipping batch.", + vectors.Count, batch.Count); + return 0; + } + + for (var i = 0; i < batch.Count; i++) + { + if (vectors[i].Length == 0) continue; // provider soft-failure for one item + batch[i].Embedding = new Pgvector.Vector(vectors[i]); + } + await db.SaveChangesAsync(ct); + _logger.LogDebug("Embedded {Count} emails", batch.Count); + return batch.Count; + } + + private static string Truncate(string s, int max) => s.Length <= max ? s : s[..max]; +} diff --git a/src/InboxIntel.Infrastructure/DependencyInjection.cs b/src/InboxIntel.Infrastructure/DependencyInjection.cs index 51fe403..9bd0be3 100644 --- a/src/InboxIntel.Infrastructure/DependencyInjection.cs +++ b/src/InboxIntel.Infrastructure/DependencyInjection.cs @@ -92,6 +92,9 @@ public static class DependencyInjection break; } services.AddScoped(); + // Semantic search: fills Email.Embedding in the background; no-ops when the + // embedding provider is unavailable (AI disabled), so lexical search is unaffected. + services.AddHostedService(); // Background worker (daily incremental sync + aggregate refresh) services.AddHostedService(); diff --git a/src/InboxIntel.Infrastructure/InboxIntel.Infrastructure.csproj b/src/InboxIntel.Infrastructure/InboxIntel.Infrastructure.csproj index 402a50a..fe5346e 100644 --- a/src/InboxIntel.Infrastructure/InboxIntel.Infrastructure.csproj +++ b/src/InboxIntel.Infrastructure/InboxIntel.Infrastructure.csproj @@ -32,4 +32,7 @@ + + + diff --git a/tests/InboxIntel.IntegrationTests/EmbeddingBackfillTests.cs b/tests/InboxIntel.IntegrationTests/EmbeddingBackfillTests.cs new file mode 100644 index 0000000..9a80fe4 --- /dev/null +++ b/tests/InboxIntel.IntegrationTests/EmbeddingBackfillTests.cs @@ -0,0 +1,95 @@ +using FluentAssertions; +using InboxIntel.Application.Abstractions; +using InboxIntel.Domain.Entities; +using InboxIntel.Infrastructure.Ai; +using InboxIntel.Infrastructure.Persistence; +using Microsoft.EntityFrameworkCore; +using Microsoft.Extensions.DependencyInjection; +using Microsoft.Extensions.Logging.Abstractions; +using Pgvector.EntityFrameworkCore; +using Xunit; + +namespace InboxIntel.IntegrationTests; + +/// +/// Semantic-search backfill: the worker embeds emails lacking an Embedding and persists the +/// vectors. Runs against live Postgres (pgvector) since the Embedding column is ignored under +/// the InMemory provider; the CI db-tests job provides the database. Uses a deterministic fake +/// embedding provider — real-Ollama integration is verified separately (endpoint contract: +/// /api/embeddings, 768 dims). +/// +[Trait("Category", "LiveDb")] +[Collection("LiveDb")] // serialise LiveDb classes: concurrent MigrateAsync on a fresh DB races +public class EmbeddingBackfillTests +{ + private static string? Conn => Environment.GetEnvironmentVariable("LIVEDB_CONNECTION"); + + private sealed class FakeCurrentUser : ICurrentUser + { + public Guid UserId => Guid.Empty; // worker scope sees all rows + public bool IsAuthenticated => false; + } + + private sealed class FakeEmbeddings : IEmbeddingProvider + { + public bool IsAvailable => true; + public Task EmbedAsync(string text, CancellationToken ct = default) + => Task.FromResult(Vec(text)); + public Task> EmbedBatchAsync(IReadOnlyList texts, CancellationToken ct = default) + => Task.FromResult>(texts.Select(Vec).ToList()); + private static float[] Vec(string text) + { + var v = new float[768]; + v[0] = text.Length; // deterministic, content-dependent + return v; + } + } + + [Fact] + public async Task Worker_embeds_pending_emails_and_persists_vectors() + { + if (Conn is null) return; // soft-skip outside the live-db CI job + var uid = Guid.NewGuid(); + var opts = new DbContextOptionsBuilder().UseNpgsql(Conn!, o => o.UseVector()).Options; + + var services = new ServiceCollection(); + services.AddScoped(); + services.AddScoped(_ => new AppDbContext(opts, new FakeCurrentUser())); + services.AddScoped(); + using var sp = services.BuildServiceProvider(); + + using (var seed = new AppDbContext(opts, new FakeCurrentUser())) + { + await seed.Database.MigrateAsync(); + seed.Users.Add(new User { Id = uid, GoogleSubjectId = "g" + uid, Email = uid + "@t.t" }); + var dom = new MailDomain { UserId = uid, Name = "t.t" }; + var snd = new Sender { UserId = uid, Address = "a@t.t", Domain = dom }; + var thr = new MailThread { UserId = uid, GmailThreadId = "th" + uid }; + seed.AddRange(dom, snd, thr); + seed.Emails.Add(new Email { UserId = uid, GmailMessageId = "e1" + uid, Subject = "hello world", Sender = snd, Thread = thr, SentAtUtc = DateTimeOffset.UtcNow }); + seed.Emails.Add(new Email { UserId = uid, GmailMessageId = "e2" + uid, Subject = "quarterly invoice", Sender = snd, Thread = thr, SentAtUtc = DateTimeOffset.UtcNow }); + await seed.SaveChangesAsync(); + } + try + { + var worker = new EmbeddingBackfillWorker( + sp.GetRequiredService(), NullLogger.Instance); + var processed = await worker.ProcessBatchAsync(CancellationToken.None); + processed.Should().BeGreaterThanOrEqualTo(2); + + using var check = new AppDbContext(opts, new FakeCurrentUser()); + var mine = await check.Emails.Where(e => e.UserId == uid).ToListAsync(); + mine.Should().OnlyContain(e => e.Embedding != null); + mine.First().Embedding!.ToArray().Length.Should().Be(768); + } + finally + { + using var c = new AppDbContext(opts, new FakeCurrentUser()); + await c.Emails.Where(e => e.UserId == uid).ExecuteDeleteAsync(); + await c.Threads.Where(t => t.UserId == uid).ExecuteDeleteAsync(); + await c.Senders.Where(s => s.UserId == uid).ExecuteDeleteAsync(); + await c.Domains.Where(d => d.UserId == uid).ExecuteDeleteAsync(); + await c.Users.Where(u => u.Id == uid).ExecuteDeleteAsync(); + } + } +} diff --git a/tests/InboxIntel.IntegrationTests/LiveDbSearchTests.cs b/tests/InboxIntel.IntegrationTests/LiveDbSearchTests.cs index 5c7fc9d..4793b15 100644 --- a/tests/InboxIntel.IntegrationTests/LiveDbSearchTests.cs +++ b/tests/InboxIntel.IntegrationTests/LiveDbSearchTests.cs @@ -21,6 +21,7 @@ namespace InboxIntel.IntegrationTests; /// InMemory test run is unaffected. /// [Trait("Category", "LiveDb")] +[Collection("LiveDb")] // serialise LiveDb classes: concurrent MigrateAsync on a fresh DB races public class LiveDbSearchTests { private static string? Conn => Environment.GetEnvironmentVariable("LIVEDB_CONNECTION");