Compare commits

..

1 Commits

Author SHA1 Message Date
cesnimda 3b93893618 feat(ai): semantic-search backfill worker + Ollama compose profile (RECOMMENDATIONS #4a)
CI / backend (pull_request) Successful in 49s
CI / frontend (pull_request) Successful in 12s
CI / format (pull_request) Successful in 50s
CI / db-tests (pull_request) Successful in 51s
Security / secrets (pull_request) Successful in 4s
Security / dependencies (pull_request) Successful in 52s
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 <noreply@anthropic.com>
2026-07-02 16:50:27 +02:00
9 changed files with 227 additions and 40 deletions
-3
View File
@@ -15,6 +15,3 @@ FRONTEND_ORIGIN=http://localhost:8081
# Set DEV_MODE=true and MAX_MESSAGES=1000 to test against a large mailbox. # Set DEV_MODE=true and MAX_MESSAGES=1000 to test against a large mailbox.
DEV_MODE=false DEV_MODE=false
MAX_MESSAGES=0 MAX_MESSAGES=0
# Nightly DB backup rotation (days of dumps to keep in ./backups)
BACKUP_KEEP_DAYS=7
-3
View File
@@ -21,9 +21,6 @@ frontend/.vite/
appsettings.*.local.json appsettings.*.local.json
secrets.json secrets.json
## DB backups (never commit dumps)
backups/
## Logs ## Logs
logs/ logs/
*.log *.log
+1 -3
View File
@@ -29,9 +29,7 @@ reverse proxy.
mitigate (a third party reading the DB files) reduces to "someone with access to your mitigate (a third party reading the DB files) reduces to "someone with access to your
machine" — mitigate it at the layer that actually works: machine" — mitigate it at the layer that actually works:
- **Use full-disk or volume encryption** on the host (BitLocker/LUKS) — strongly recommended. - **Use full-disk or volume encryption** on the host (BitLocker/LUKS) — strongly recommended.
- **Encrypt backups**: nightly `pg_dump` rotation runs via the compose `backup` service - **Encrypt backups** of the `pgdata` volume the same way.
into `./backups/` (git-ignored) — keep that directory on an encrypted disk and copy it
off-machine. Restore: `docker compose exec -T postgres psql -U inboxintel -d inboxintel < backups/<file>.sql`.
- Before any **multi-user** deployment, revisit per the multi-provider security design - Before any **multi-user** deployment, revisit per the multi-provider security design
(host admins must not be able to read members' mail — plaintext bodies break that promise). (host admins must not be able to read members' mail — plaintext bodies break that promise).
2. **DB connection is not TLS** — Postgres is only reachable on the compose-internal network / 2. **DB connection is not TLS** — Postgres is only reachable on the compose-internal network /
+20 -31
View File
@@ -22,37 +22,6 @@ services:
timeout: 5s timeout: 5s
retries: 10 retries: 10
# Nightly logical backups (RECOMMENDATIONS #3 — previously there were NONE). Dumps
# rotate after BACKUP_KEEP_DAYS. The ./backups host directory should live on an
# encrypted disk and be included in your off-machine backup regime (see SECURITY.md).
# Restore: docker compose exec -T postgres psql -U inboxintel -d inboxintel < backups/<file>.sql
backup:
image: pgvector/pgvector:pg16
entrypoint: /bin/sh
command:
- -c
- |
while true; do
ts=$$(date -u +%Y%m%d-%H%M%S)
if pg_dump -h postgres -U inboxintel -d inboxintel > /backups/inboxintel-$$ts.sql.tmp; then
mv /backups/inboxintel-$$ts.sql.tmp /backups/inboxintel-$$ts.sql
echo "backup OK: inboxintel-$$ts.sql"
else
rm -f /backups/inboxintel-$$ts.sql.tmp
echo "backup FAILED at $$ts" >&2
fi
find /backups -name 'inboxintel-*.sql' -mtime +$${BACKUP_KEEP_DAYS:-7} -delete
sleep 86400
done
environment:
PGPASSWORD: ${POSTGRES_PASSWORD:?set POSTGRES_PASSWORD in deploy/.env}
BACKUP_KEEP_DAYS: ${BACKUP_KEEP_DAYS:-7}
volumes:
- ./backups:/backups
depends_on:
postgres:
condition: service_healthy
api: api:
build: build:
context: . context: .
@@ -65,6 +34,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 +61,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 +93,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");