feat(ai): semantic-search backfill + Ollama profile (#30)
CI / backend (push) Successful in 55s
CI / frontend (push) Successful in 12s
CI / format (push) Successful in 48s
CI / db-tests (push) Successful in 55s
Deploy Staging / deploy (push) Successful in 42s
CI / backend (pull_request) Successful in 54s
CI / frontend (pull_request) Successful in 13s
CI / format (pull_request) Successful in 49s
CI / db-tests (pull_request) Successful in 56s
Security / secrets (push) Successful in 4s
Security / dependencies (push) Successful in 58s
Security / secrets (pull_request) Successful in 4s
Security / dependencies (pull_request) Successful in 55s
CI / backend (push) Successful in 55s
CI / frontend (push) Successful in 12s
CI / format (push) Successful in 48s
CI / db-tests (push) Successful in 55s
Deploy Staging / deploy (push) Successful in 42s
CI / backend (pull_request) Successful in 54s
CI / frontend (pull_request) Successful in 13s
CI / format (pull_request) Successful in 49s
CI / db-tests (pull_request) Successful in 56s
Security / secrets (push) Successful in 4s
Security / dependencies (push) Successful in 58s
Security / secrets (pull_request) Successful in 4s
Security / dependencies (pull_request) Successful in 55s
This commit was merged in pull request #30.
This commit is contained in:
@@ -65,6 +65,8 @@ services:
|
|||||||
GoogleOAuth__ClientId: ${GOOGLE_CLIENT_ID:-}
|
GoogleOAuth__ClientId: ${GOOGLE_CLIENT_ID:-}
|
||||||
GoogleOAuth__ClientSecret: ${GOOGLE_CLIENT_SECRET:-}
|
GoogleOAuth__ClientSecret: ${GOOGLE_CLIENT_SECRET:-}
|
||||||
Ai__Mode: ${AI_MODE:-Disabled}
|
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
|
# 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.
|
# and MAX_MESSAGES=1000 in deploy/.env to exercise it in this Docker setup.
|
||||||
App__DevMode: ${DEV_MODE:-false}
|
App__DevMode: ${DEV_MODE:-false}
|
||||||
@@ -90,6 +92,23 @@ services:
|
|||||||
ports:
|
ports:
|
||||||
- "8081:80"
|
- "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
|
# Optional reverse proxy. Enable with: docker compose --profile proxy up
|
||||||
nginx:
|
nginx:
|
||||||
image: nginx:alpine
|
image: nginx:alpine
|
||||||
@@ -105,3 +124,4 @@ services:
|
|||||||
volumes:
|
volumes:
|
||||||
pgdata:
|
pgdata:
|
||||||
keys:
|
keys:
|
||||||
|
ollama:
|
||||||
|
|||||||
@@ -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;
|
||||||
|
|
||||||
|
/// <summary>
|
||||||
|
/// Fills <c>Email.Embedding</c> (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.
|
||||||
|
/// </summary>
|
||||||
|
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<EmbeddingBackfillWorker> _logger;
|
||||||
|
|
||||||
|
public EmbeddingBackfillWorker(IServiceScopeFactory scopeFactory, ILogger<EmbeddingBackfillWorker> 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<IEmbeddingProvider>().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);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/// <summary>Embeds one batch. Public-ish (internal) for direct testing.</summary>
|
||||||
|
internal async Task<int> ProcessBatchAsync(CancellationToken ct)
|
||||||
|
{
|
||||||
|
using var scope = _scopeFactory.CreateScope();
|
||||||
|
var db = scope.ServiceProvider.GetRequiredService<AppDbContext>();
|
||||||
|
var embeddings = scope.ServiceProvider.GetRequiredService<IEmbeddingProvider>();
|
||||||
|
|
||||||
|
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];
|
||||||
|
}
|
||||||
@@ -92,6 +92,9 @@ public static class DependencyInjection
|
|||||||
break;
|
break;
|
||||||
}
|
}
|
||||||
services.AddScoped<IAiService, AiService>();
|
services.AddScoped<IAiService, AiService>();
|
||||||
|
// 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<EmbeddingBackfillWorker>();
|
||||||
|
|
||||||
// Background worker (daily incremental sync + aggregate refresh)
|
// Background worker (daily incremental sync + aggregate refresh)
|
||||||
services.AddHostedService<GmailSyncWorker>();
|
services.AddHostedService<GmailSyncWorker>();
|
||||||
|
|||||||
@@ -32,4 +32,7 @@
|
|||||||
<ProjectReference Include="..\InboxIntel.Application\InboxIntel.Application.csproj" />
|
<ProjectReference Include="..\InboxIntel.Application\InboxIntel.Application.csproj" />
|
||||||
<ProjectReference Include="..\InboxIntel.Domain\InboxIntel.Domain.csproj" />
|
<ProjectReference Include="..\InboxIntel.Domain\InboxIntel.Domain.csproj" />
|
||||||
</ItemGroup>
|
</ItemGroup>
|
||||||
|
<ItemGroup>
|
||||||
|
<InternalsVisibleTo Include="InboxIntel.IntegrationTests" />
|
||||||
|
</ItemGroup>
|
||||||
</Project>
|
</Project>
|
||||||
|
|||||||
@@ -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;
|
||||||
|
|
||||||
|
/// <summary>
|
||||||
|
/// 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).
|
||||||
|
/// </summary>
|
||||||
|
[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<float[]> EmbedAsync(string text, CancellationToken ct = default)
|
||||||
|
=> Task.FromResult(Vec(text));
|
||||||
|
public Task<IReadOnlyList<float[]>> EmbedBatchAsync(IReadOnlyList<string> texts, CancellationToken ct = default)
|
||||||
|
=> Task.FromResult<IReadOnlyList<float[]>>(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<AppDbContext>().UseNpgsql(Conn!, o => o.UseVector()).Options;
|
||||||
|
|
||||||
|
var services = new ServiceCollection();
|
||||||
|
services.AddScoped<ICurrentUser, FakeCurrentUser>();
|
||||||
|
services.AddScoped(_ => new AppDbContext(opts, new FakeCurrentUser()));
|
||||||
|
services.AddScoped<IEmbeddingProvider, FakeEmbeddings>();
|
||||||
|
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<IServiceScopeFactory>(), NullLogger<EmbeddingBackfillWorker>.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();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -21,6 +21,7 @@ namespace InboxIntel.IntegrationTests;
|
|||||||
/// InMemory test run is unaffected.
|
/// InMemory test run is unaffected.
|
||||||
/// </summary>
|
/// </summary>
|
||||||
[Trait("Category", "LiveDb")]
|
[Trait("Category", "LiveDb")]
|
||||||
|
[Collection("LiveDb")] // serialise LiveDb classes: concurrent MigrateAsync on a fresh DB races
|
||||||
public class LiveDbSearchTests
|
public class LiveDbSearchTests
|
||||||
{
|
{
|
||||||
private static string? Conn => Environment.GetEnvironmentVariable("LIVEDB_CONNECTION");
|
private static string? Conn => Environment.GetEnvironmentVariable("LIVEDB_CONNECTION");
|
||||||
|
|||||||
Reference in New Issue
Block a user