feat(worker): S3 3b härda UPDATE_USER_MEMORY — budget, kostnad, anti-påhitt, GDPR-regression
This commit is contained in:
@@ -19,6 +19,8 @@ export interface WorkerContext {
|
||||
db: Database;
|
||||
aamos: AamosClient;
|
||||
redis?: Redis;
|
||||
/** Delad dagsbudget-store så att memory-jobbet kan spegla skanningens budgetkoll. */
|
||||
budgetStore?: BudgetStore;
|
||||
/** Bygger läs-URL för lagrade bilder (samma signaturlogik som API:ts mock-S3). */
|
||||
readUrl: (key: string) => string;
|
||||
apiBaseUrl: string;
|
||||
@@ -70,6 +72,7 @@ export function createContext(redis?: Redis): WorkerContext {
|
||||
db,
|
||||
aamos,
|
||||
redis,
|
||||
budgetStore,
|
||||
apiBaseUrl,
|
||||
readUrl: (key: string) => {
|
||||
const sig = createHmac("sha256", secret).update(key).digest("hex").slice(0, 32);
|
||||
|
||||
@@ -0,0 +1,41 @@
|
||||
import { sql } from "drizzle-orm";
|
||||
import { schema, type Database } from "@app/database";
|
||||
|
||||
function currentMonth(): string {
|
||||
const d = new Date();
|
||||
return `${d.getUTCFullYear()}-${String(d.getUTCMonth() + 1).padStart(2, "0")}`;
|
||||
}
|
||||
|
||||
/**
|
||||
* Bokför AI-användning utan PII i ai_usage_counters.
|
||||
* Används av både scan-processor och UPDATE_USER_MEMORY.
|
||||
*/
|
||||
export async function recordAiUsage(
|
||||
db: Database,
|
||||
userId: string,
|
||||
usage: { costUsd: number; inputTokens: number; outputTokens: number; aiScans?: number },
|
||||
): Promise<void> {
|
||||
const microcents = Math.round(usage.costUsd * 100_000_000);
|
||||
const month = currentMonth();
|
||||
const aiScans = usage.aiScans ?? 0;
|
||||
|
||||
await db
|
||||
.insert(schema.aiUsageCounters)
|
||||
.values({
|
||||
userId,
|
||||
month,
|
||||
aiScans,
|
||||
aiTokensIn: usage.inputTokens,
|
||||
aiTokensOut: usage.outputTokens,
|
||||
aiCostUsdMicrocents: microcents,
|
||||
})
|
||||
.onConflictDoUpdate({
|
||||
target: [schema.aiUsageCounters.userId, schema.aiUsageCounters.month],
|
||||
set: {
|
||||
aiScans: sql`${schema.aiUsageCounters.aiScans} + ${aiScans}`,
|
||||
aiTokensIn: sql`${schema.aiUsageCounters.aiTokensIn} + ${usage.inputTokens}`,
|
||||
aiTokensOut: sql`${schema.aiUsageCounters.aiTokensOut} + ${usage.outputTokens}`,
|
||||
aiCostUsdMicrocents: sql`${schema.aiUsageCounters.aiCostUsdMicrocents} + ${microcents}`,
|
||||
},
|
||||
});
|
||||
}
|
||||
@@ -5,6 +5,12 @@ import { deriveMemoryUpdates } from "@app/memory-client";
|
||||
import { getLocaleContext } from "../locale.js";
|
||||
import type { WorkerContext } from "../context.js";
|
||||
import { cancelTimedOutCookingSessions } from "@app/database";
|
||||
import { recordAiUsage } from "../lib/ai-usage.js";
|
||||
|
||||
/** Pessimistisk kostnadsuppskattning för UPDATE_USER_MEMORY (text-only). */
|
||||
const COST_ESTIMATE_PER_MEMORY_SYNC_USD = 0.001;
|
||||
/** Maxkonfidens för ai_inferred-minnen (R1: låg startkonfidens). */
|
||||
const AI_INFERRED_MAX_CONFIDENCE = 0.7;
|
||||
|
||||
export { processProactiveTips } from "./proactive-tips.js";
|
||||
|
||||
@@ -237,6 +243,8 @@ export async function processMealBoxReminders(ctx: WorkerContext): Promise<numbe
|
||||
* förslag skrivs till memory_items där användaren äger dem.
|
||||
*/
|
||||
export async function processMemorySync(ctx: WorkerContext): Promise<number> {
|
||||
const dailyBudgetUsd = Number(process.env.GEMINI_DAILY_BUDGET_USD ?? 0);
|
||||
|
||||
const users = await ctx.db
|
||||
.select({ userId: schema.userConsents.userId })
|
||||
.from(schema.userConsents)
|
||||
@@ -250,7 +258,12 @@ export async function processMemorySync(ctx: WorkerContext): Promise<number> {
|
||||
let updates = 0;
|
||||
for (const { userId } of users) {
|
||||
const events = await ctx.db
|
||||
.select()
|
||||
.select({
|
||||
id: schema.domainEvents.id,
|
||||
type: schema.domainEvents.type,
|
||||
occurredAt: schema.domainEvents.occurredAt,
|
||||
payload: schema.domainEvents.payload,
|
||||
})
|
||||
.from(schema.domainEvents)
|
||||
.where(
|
||||
and(
|
||||
@@ -267,6 +280,16 @@ export async function processMemorySync(ctx: WorkerContext): Promise<number> {
|
||||
.limit(100);
|
||||
if (events.length < 3) continue;
|
||||
|
||||
// BUDGET: spegla skanningen — över budget får inget AAMOS-anrop göras.
|
||||
if (dailyBudgetUsd > 0 && ctx.budgetStore) {
|
||||
const currentSpend = await ctx.budgetStore.getDailySpendUsd();
|
||||
if (currentSpend + COST_ESTIMATE_PER_MEMORY_SYNC_USD > dailyBudgetUsd) {
|
||||
// eslint-disable-next-line no-console
|
||||
console.log(`UPDATE_USER_MEMORY: skippar ${userId}: daglig budget förbrukad.`);
|
||||
continue;
|
||||
}
|
||||
}
|
||||
|
||||
const existing = await ctx.db
|
||||
.select({ key: schema.memoryItems.key })
|
||||
.from(schema.memoryItems)
|
||||
@@ -279,11 +302,12 @@ export async function processMemorySync(ctx: WorkerContext): Promise<number> {
|
||||
const has = (kind: string) => consents.find((c) => c.kind === kind)?.status === "granted";
|
||||
|
||||
const localeContext = await getLocaleContext(ctx, userId);
|
||||
const proposals = await deriveMemoryUpdates(ctx.aamos, {
|
||||
const { proposals, usage } = await deriveMemoryUpdates(ctx.aamos, {
|
||||
scope: "user",
|
||||
scopeId: userId,
|
||||
localeContext,
|
||||
events: events.map((e) => ({
|
||||
id: e.id,
|
||||
type: e.type,
|
||||
occurredAt: e.occurredAt.toISOString(),
|
||||
payload: e.payload,
|
||||
@@ -296,7 +320,35 @@ export async function processMemorySync(ctx: WorkerContext): Promise<number> {
|
||||
},
|
||||
});
|
||||
|
||||
// Bokför AI-kostnad/tokens så total dagskostnad inkluderar minnes-AI:n.
|
||||
if (usage) {
|
||||
await recordAiUsage(ctx.db, userId, {
|
||||
costUsd: usage.costUsd,
|
||||
inputTokens: usage.inputTokens,
|
||||
outputTokens: usage.outputTokens,
|
||||
});
|
||||
}
|
||||
|
||||
const eventIdSet = new Set(events.map((e) => e.id));
|
||||
|
||||
for (const proposal of proposals) {
|
||||
// R1: varje förslag måste vara grundat i faktiska events.
|
||||
const hasEventSupport =
|
||||
proposal.sourceEventIds.length > 0 &&
|
||||
proposal.sourceEventIds.every((id) => eventIdSet.has(id));
|
||||
if (!hasEventSupport) {
|
||||
// eslint-disable-next-line no-console
|
||||
console.log(`UPDATE_USER_MEMORY: avvisar ${proposal.key} för ${userId}: saknar event-stöd.`);
|
||||
continue;
|
||||
}
|
||||
|
||||
// R1: ai_inferred ska ha låg konfidens.
|
||||
if (proposal.origin === "ai_inferred" && proposal.confidence > AI_INFERRED_MAX_CONFIDENCE) {
|
||||
// eslint-disable-next-line no-console
|
||||
console.log(`UPDATE_USER_MEMORY: avvisar ${proposal.key} för ${userId}: ai_inferred med för hög confidence.`);
|
||||
continue;
|
||||
}
|
||||
|
||||
// Skriv aldrig över användarverifierade poster (spec §30: user_stated vinner).
|
||||
const [current] = await ctx.db
|
||||
.select()
|
||||
|
||||
@@ -1,14 +1,11 @@
|
||||
import { eq, sql } from "drizzle-orm";
|
||||
import { eq } from "drizzle-orm";
|
||||
import { schema, trackProductAnalytics } from "@app/database";
|
||||
import { scanCompleted, scanFailed } from "@app/analytics";
|
||||
import type { AamosResult, AamosTaskType, DetectedItem } from "@app/ai-contracts";
|
||||
import type { WorkerContext } from "../context.js";
|
||||
import { getLocaleContext } from "../locale.js";
|
||||
import type { LocaleContext } from "@app/shared-types";
|
||||
function currentMonth(): string {
|
||||
const d = new Date();
|
||||
return `${d.getUTCFullYear()}-${String(d.getUTCMonth() + 1).padStart(2, "0")}`;
|
||||
}
|
||||
import { recordAiUsage } from "../lib/ai-usage.js";
|
||||
|
||||
/**
|
||||
* Bild-/OCR-jobb (spec §54): hämtar scan_job, anropar AAMOS med kontraktvaliderad
|
||||
@@ -84,7 +81,14 @@ export async function processScanJob(ctx: WorkerContext, scanJobId: string): Pro
|
||||
await recordScanCompleted(ctx, job, result);
|
||||
|
||||
// Bokför verklig AI-kostnad/tokens utan PII (spec §45).
|
||||
await recordAiUsage(ctx, job.userId, result);
|
||||
if (result.costUsd != null || result.inputTokens || result.outputTokens) {
|
||||
await recordAiUsage(ctx.db, job.userId, {
|
||||
costUsd: result.costUsd ?? 0,
|
||||
inputTokens: result.inputTokens ?? 0,
|
||||
outputTokens: result.outputTokens ?? 0,
|
||||
aiScans: 1,
|
||||
});
|
||||
}
|
||||
|
||||
// MEAL_PHOTO_ANALYZED-event för tallriksfoton (spec §55)
|
||||
if (job.jobType === "ANALYZE_MEAL_IMAGE") {
|
||||
@@ -294,41 +298,6 @@ function tokenize(text: string): string[] {
|
||||
.filter((t) => t.length > 1);
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
// AI-kostnadsbokföring utan PII.
|
||||
// ---------------------------------------------------------------------------
|
||||
|
||||
async function recordAiUsage(
|
||||
ctx: WorkerContext,
|
||||
userId: string,
|
||||
result: AamosResult<AamosTaskType>,
|
||||
): Promise<void> {
|
||||
const costUsd = result.costUsd ?? 0;
|
||||
const tokensIn = result.inputTokens ?? 0;
|
||||
const tokensOut = result.outputTokens ?? 0;
|
||||
const microcents = Math.round(costUsd * 100_000_000);
|
||||
const month = currentMonth();
|
||||
|
||||
await ctx.db
|
||||
.insert(schema.aiUsageCounters)
|
||||
.values({
|
||||
userId,
|
||||
month,
|
||||
aiScans: 1,
|
||||
aiTokensIn: tokensIn,
|
||||
aiTokensOut: tokensOut,
|
||||
aiCostUsdMicrocents: microcents,
|
||||
})
|
||||
.onConflictDoUpdate({
|
||||
target: [schema.aiUsageCounters.userId, schema.aiUsageCounters.month],
|
||||
set: {
|
||||
aiScans: sql`${schema.aiUsageCounters.aiScans} + 1`,
|
||||
aiTokensIn: sql`${schema.aiUsageCounters.aiTokensIn} + ${tokensIn}`,
|
||||
aiTokensOut: sql`${schema.aiUsageCounters.aiTokensOut} + ${tokensOut}`,
|
||||
aiCostUsdMicrocents: sql`${schema.aiUsageCounters.aiCostUsdMicrocents} + ${microcents}`,
|
||||
},
|
||||
});
|
||||
}
|
||||
|
||||
function classifyScanError(error?: string | null): string {
|
||||
if (!error) return "unknown";
|
||||
|
||||
Reference in New Issue
Block a user