Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 7 additions & 0 deletions .env.example
Original file line number Diff line number Diff line change
Expand Up @@ -130,6 +130,13 @@ MCP_RATE_LIMIT_PER_MINUTE=60
# KG_LLM_MODEL=gpt-4o-mini
# OPENAI_API_KEY= # or OPENROUTER_API_KEY / ANTHROPIC_API_KEY
#
# Local / OpenAI-compatible endpoint (Ollama, LM Studio, ...). When set, the KG
# LLM path talks to this base URL instead of OpenAI/OpenRouter/Anthropic, still
# using KG_LLM_MODEL. No key needed — leave KG_LLM_API_KEY empty; it is only
# sent as an Authorization header when set. Validated against local qwen2.5.
# KG_LLM_BASE_URL=http://localhost:11434/v1
# KG_LLM_API_KEY=
#
# Scheduled AI extension (cloud cron): periodically extends the graph + skills
# from captured user intents. Off unless ALL of: this flag is true, the cron
# runs (POST /api/cron/kg-discovery with CRON_SECRET), and the workspace turned
Expand Down
14 changes: 14 additions & 0 deletions docker-compose.yml
Original file line number Diff line number Diff line change
Expand Up @@ -81,6 +81,20 @@ services:
# address. Needed for the bundled MOTIS (`motis`) and any other API you
# run on the same Docker network.
- SSRF_ALLOWED_HOSTS=${SSRF_ALLOWED_HOSTS:-}
# Knowledge Graph LLM (optional): hosted providers need their API key,
# a local OpenAI-compatible endpoint needs KG_LLM_BASE_URL and no key
# (KG_LLM_API_KEY is only sent when set). From inside this container
# `localhost` is the app itself — point at the host or a sibling
# service instead, e.g. http://host.docker.internal:11434/v1
# (Linux: extra_hosts: ["host.docker.internal:host-gateway"]).
- KG_LLM_ENABLED=${KG_LLM_ENABLED:-}
- KG_LLM_PROVIDER=${KG_LLM_PROVIDER:-}
- KG_LLM_MODEL=${KG_LLM_MODEL:-}
- KG_LLM_BASE_URL=${KG_LLM_BASE_URL:-}
- KG_LLM_API_KEY=${KG_LLM_API_KEY:-}
- OPENAI_API_KEY=${OPENAI_API_KEY:-}
- OPENROUTER_API_KEY=${OPENROUTER_API_KEY:-}
- ANTHROPIC_API_KEY=${ANTHROPIC_API_KEY:-}

depends_on:
postgres:
Expand Down
19 changes: 19 additions & 0 deletions docs/knowledge-graph.md
Original file line number Diff line number Diff line change
Expand Up @@ -179,6 +179,16 @@ KG_LLM_PROVIDER=openai # openai | openrouter | anthropic
KG_LLM_MODEL=gpt-4o-mini # anthropic default: claude-haiku-4-5
OPENAI_API_KEY= # or OPENROUTER_API_KEY / ANTHROPIC_API_KEY

# Local / OpenAI-compatible endpoint (Ollama, LM Studio, ...). When set, the
# KG LLM path talks here instead of the hosted providers, still using
# KG_LLM_MODEL. No key needed: KG_LLM_API_KEY is optional and only sent as an
# Authorization header when set. No response_format is sent on this path, so
# the model answers from the JSON instruction; unparsable replies are logged
# (model + endpoint origin) and that pass is skipped. Validated against a local
# qwen2.5 (Qwen2.5-0.5B-Instruct via an OpenAI-compatible front).
# KG_LLM_BASE_URL=http://localhost:11434/v1
# KG_LLM_API_KEY=

# Scheduled extension (cloud cron) — all must align: this flag, the cron call,
# and the per-workspace "Scheduled AI extension" switch
KG_LLM_CRON_ENABLED=false
Expand All @@ -190,6 +200,15 @@ KG_LLM_BATCH=false # Anthropic Message Batches (~50% cheaper)
KG_LLM_REDACT_INTENTS=true
```

Docker note: `localhost` inside the app container means the AnythingMCP
container itself, not your host. To reach an Ollama running on the host, use
`http://host.docker.internal:11434/v1` — on Linux that name needs an explicit
mapping (`extra_hosts: ["host.docker.internal:host-gateway"]`). When Ollama
runs as a service in the same Compose project/network, use its service name
instead (e.g. `http://ollama:11434/v1`). The `KG_LLM_*` variables and the
provider keys are passed through to the app container in `docker-compose.yml`
with empty defaults, so setting them in `.env` is enough.

The graph itself (static + observational layers, manual editing, the
`kg_how_to_obtain` tool) works with **no LLM key** — the AI flags only add the
optional enrichment and skill-generation passes on top.
72 changes: 72 additions & 0 deletions packages/backend/src/knowledge-graph/kg-llm.service.spec.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,72 @@
import { KgLlmService } from './kg-llm.service';

const ORG = 'org-A';

function setKgEnv() {
process.env.KG_LLM_ENABLED = 'true';
process.env.KG_LLM_BASE_URL = 'http://localhost:11434/v1';
process.env.KG_LLM_MODEL = 'qwen2.5';
delete process.env.KG_LLM_API_KEY;
}

function clearKgEnv() {
delete process.env.KG_LLM_ENABLED;
delete process.env.KG_LLM_BASE_URL;
delete process.env.KG_LLM_MODEL;
}

function okJson(body: any) {
return { ok: true, status: 200, text: async () => JSON.stringify(body), json: async () => body };
}

function make(prisma: any) {
const kgStatic = { isEnabled: async () => true, getFlag: async () => true };
return new KgLlmService(prisma, kgStatic as any);
}

function nodes() {
return [0, 1].map((i) => ({
id: `n${i}`,
entity: `ent${i}`,
fields: [{ name: 'email' }],
outputFields: [],
connector: { name: 'crm' },
}));
}

/** Unusable custom reply must not store the hash (test 1 for #816 review). */
describe('KgLlmService.enrich custom skip', () => {
const realFetch = global.fetch;

beforeEach(() => {
setKgEnv();
global.fetch = jest.fn();
});
afterEach(() => {
global.fetch = realFetch;
clearKgEnv();
jest.restoreAllMocks();
});

it('does not update kg_llm_hash on an unusable reply and stays retryable', async () => {
(global.fetch as jest.Mock).mockResolvedValue(
okJson({ choices: [{ message: { content: 'Sure, here are some thoughts...' } }] }),
);
const prisma = {
kgNode: { findMany: jest.fn().mockResolvedValue(nodes()) },
orgSettings: { findUnique: jest.fn().mockResolvedValue(null), upsert: jest.fn() },
};
const res = await make(prisma).enrich(ORG);
expect(res).toEqual({ suggested: 0, skipped: true, model: 'qwen2.5' });
expect(prisma.orgSettings.upsert).not.toHaveBeenCalled();

// Retry with a usable reply stores the hash normally.
(global.fetch as jest.Mock).mockResolvedValue(
okJson({ choices: [{ message: { content: '{"relationships":[]}' } }] }),
);
const retry = await make(prisma).enrich(ORG);
expect(retry.suggested).toBe(0);
expect(retry.skipped).toBeUndefined();
expect(prisma.orgSettings.upsert).toHaveBeenCalledTimes(1);
});
});
11 changes: 8 additions & 3 deletions packages/backend/src/knowledge-graph/kg-llm.service.ts
Original file line number Diff line number Diff line change
Expand Up @@ -27,7 +27,8 @@ Rules:

/**
* Optional LLM-assisted KG enrichment. Opt-in twice: a global env flag
* (KG_LLM_ENABLED + an API key) AND a per-workspace switch (kg_llm_enabled,
* (KG_LLM_ENABLED + a resolvable LLM config — an API key, or KG_LLM_BASE_URL
* for local models which need no key) AND a per-workspace switch (kg_llm_enabled,
* default off, because it costs money). PII-safe: only entity + field NAMES are
* sent to the model — never values. Results are stored as suggested LLM edges
* for a human to confirm. Cached by a content hash so unchanged graphs don't
Expand All @@ -42,7 +43,7 @@ export class KgLlmService {
private readonly kgStatic: KgStaticService,
) {}

/** Globally available (env flag + a configured API key). */
/** Globally available (env flag + a resolvable LLM config). */
globallyAvailable(): boolean {
return process.env.KG_LLM_ENABLED === 'true' && !!resolveLlmConfig();
}
Expand All @@ -66,7 +67,11 @@ export class KgLlmService {
if (!built) return { suggested: 0, model: cfg.model };
if ('skipped' in built) return { suggested: 0, skipped: true, model: cfg.model };

const { json, usage } = await chatJson(cfg, built.system, built.user);
const { json, skipped, usage } = await chatJson(cfg, built.system, built.user);
// Custom endpoint with an unusable reply: return early before
// applyEnrichResult so kg_llm_hash is NOT stored and the next run retries
// normally instead of seeing an unchanged graph and reporting `skipped`.
if (skipped) return { suggested: 0, skipped: true, model: cfg.model };
const suggested = await this.applyEnrichResult(organizationId, json, built);
this.logger.log(
`KG LLM enrich ${organizationId}: ${suggested} suggested (${cfg.model}, in=${usage?.inputTokens ?? '?'} out=${usage?.outputTokens ?? '?'})`,
Expand Down
104 changes: 104 additions & 0 deletions packages/backend/src/knowledge-graph/kg-skill.service.spec.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,10 @@
import { ConflictException, NotFoundException } from '@nestjs/common';
import { KgSkillService } from './kg-skill.service';

function okJson(body: any) {
return { ok: true, status: 200, text: async () => JSON.stringify(body), json: async () => body };
}

/** Tenant isolation + defaults for manual skill creation. Prisma/LLM mocked. */
describe('KgSkillService.create', () => {
const ORG = 'org-A';
Expand Down Expand Up @@ -58,3 +62,103 @@ describe('KgSkillService.create', () => {
);
});
});

/** Unusable custom replies must leave pending suggestions untouched (test 2 for #816 review). */
describe('KgSkillService custom skip', () => {
const ORG = 'org-A';
const realFetch = global.fetch;

function setKgEnv() {
process.env.KG_LLM_BASE_URL = 'http://localhost:11434/v1';
process.env.KG_LLM_MODEL = 'qwen2.5';
delete process.env.KG_LLM_API_KEY;
}
function clearKgEnv() {
delete process.env.KG_LLM_BASE_URL;
delete process.env.KG_LLM_MODEL;
}

beforeEach(() => {
setKgEnv();
global.fetch = jest.fn();
});
afterEach(() => {
global.fetch = realFetch;
clearKgEnv();
jest.restoreAllMocks();
});

function generatePrisma() {
return {
toolInvocation: {
findMany: jest.fn().mockResolvedValue([
{
intent: 'quote net price',
status: 'SUCCESS',
tool: { name: 'get_price', connector: { name: 'Billing' } },
},
]),
},
connector: { findMany: jest.fn().mockResolvedValue([{ id: 'c1', name: 'Billing' }]) },
orgSettings: { findUnique: jest.fn().mockResolvedValue(null) },
kgSkillSuggestion: { deleteMany: jest.fn(), create: jest.fn() },
};
}

it('generate() leaves pending suggestions untouched on an unusable reply', async () => {
(global.fetch as jest.Mock).mockResolvedValue(
okJson({ choices: [{ message: { content: 'Sure, here are some thoughts...' } }] }),
);
const prisma = generatePrisma();
const svc = new KgSkillService(prisma as any, { isEnabled: async () => true } as any);
const res = await svc.generate(ORG);
expect(res).toEqual({ created: 0, skipped: true, model: 'qwen2.5' });
expect(prisma.kgSkillSuggestion.deleteMany).not.toHaveBeenCalled();
expect(prisma.kgSkillSuggestion.create).not.toHaveBeenCalled();
});

it('generate() still replaces pending suggestions on a usable reply', async () => {
(global.fetch as jest.Mock).mockResolvedValue(
okJson({
choices: [
{
message: {
content: JSON.stringify({
skills: [{ connector: 'Billing', title: 'T', instruction: 'I', confidence: 0.5 }],
}),
},
},
],
}),
);
const prisma = generatePrisma();
const svc = new KgSkillService(prisma as any, { isEnabled: async () => true } as any);
const res = await svc.generate(ORG);
expect(res.created).toBe(1);
expect(res.skipped).toBeUndefined();
expect(prisma.kgSkillSuggestion.deleteMany).toHaveBeenCalledTimes(1);
});

it('consolidate() leaves applied skills untouched on an unusable reply', async () => {
(global.fetch as jest.Mock).mockResolvedValue(
okJson({ choices: [{ message: { content: 'Sure, here are some thoughts...' } }] }),
);
const applied = [0, 1].map((i) => ({
id: `s${i}`,
title: `rule ${i}`,
whenToUse: 'w',
instruction: 'do it',
connector: { name: 'Billing' },
}));
const prisma = {
kgSkillSuggestion: { findMany: jest.fn().mockResolvedValue(applied), deleteMany: jest.fn() },
connector: { findMany: jest.fn() },
};
const svc = new KgSkillService(prisma as any, { isEnabled: async () => true } as any);
const res = await svc.consolidate(ORG);
expect(res).toEqual(
expect.objectContaining({ before: 2, after: 2, model: 'qwen2.5' }),
);
expect(prisma.kgSkillSuggestion.deleteMany).not.toHaveBeenCalled();
});
});
12 changes: 9 additions & 3 deletions packages/backend/src/knowledge-graph/kg-skill.service.ts
Original file line number Diff line number Diff line change
Expand Up @@ -63,7 +63,7 @@ export class KgSkillService {
async generate(
organizationId: string,
opts?: { mcpServerId?: string },
): Promise<{ created: number; model?: string; usage?: any }> {
): Promise<{ created: number; skipped?: boolean; model?: string; usage?: any }> {
if (!(await this.llm.isEnabled(organizationId))) {
throw new ConflictException('AI features are disabled for this workspace.');
}
Expand All @@ -85,7 +85,10 @@ export class KgSkillService {
const cfg = resolveLlmConfig()!;
const built = await this.buildConnectorRequest(organizationId);
if (!built) return { created: 0, model: cfg.model };
const { json, usage } = await chatJson(cfg, built.system, built.user);
const { json, skipped, usage } = await chatJson(cfg, built.system, built.user);
// Custom endpoint with an unusable reply: return early before
// applyConnectorResult so pending suggestions are left untouched.
if (skipped) return { created: 0, skipped: true, model: cfg.model };
const created = await this.applyConnectorResult(organizationId, json);
this.logger.log(`KG skills (connectors) ${organizationId}: ${created}`);
return { created, model: cfg.model, usage };
Expand Down Expand Up @@ -182,11 +185,14 @@ export class KgSkillService {
ok: i.status === 'SUCCESS',
}));

const { json, usage } = await chatJson(
const { json, skipped, usage } = await chatJson(
cfg,
SERVER_PROMPT,
JSON.stringify({ server: server.name, connectors: connectorsContext, calls }),
);
// Custom endpoint with an unusable reply: return early before the
// deleteMany below so pending suggestions are left untouched.
if (skipped) return { created: 0, skipped: true, model: cfg.model };
const skills: any[] = Array.isArray(json?.skills) ? json.skills : [];

await this.prisma.kgSkillSuggestion.deleteMany({
Expand Down
Loading
Loading