From a029f466176f45665d55dd7da314250067cbabc3 Mon Sep 17 00:00:00 2001 From: owk-owk130 Date: Tue, 1 Sep 2026 17:11:00 +0900 Subject: [PATCH 1/6] =?UTF-8?q?feat(server):=20knowledge=5Fsources=20?= =?UTF-8?q?=E3=83=86=E3=83=BC=E3=83=96=E3=83=AB=E3=82=92=E8=BF=BD=E5=8A=A0?= =?UTF-8?q?=E3=81=99=E3=82=8B?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-Authored-By: Claude Fable 5 Claude-Session: https://claude.ai/code/session_018BjNFvQjDzMwF8HaZHnJjP --- server/src/__tests__/helpers/test-db.ts | 18 ++++++++++++++++ .../db/migrations/0029_knowledge_sources.sql | 20 ++++++++++++++++++ server/src/db/schema.ts | 21 +++++++++++++++++++ 3 files changed, 59 insertions(+) create mode 100644 server/src/db/migrations/0029_knowledge_sources.sql diff --git a/server/src/__tests__/helpers/test-db.ts b/server/src/__tests__/helpers/test-db.ts index 99da1d07..14e14355 100644 --- a/server/src/__tests__/helpers/test-db.ts +++ b/server/src/__tests__/helpers/test-db.ts @@ -160,6 +160,24 @@ export const createTestDb = async () => { created_at TEXT NOT NULL ); + CREATE TABLE IF NOT EXISTS knowledge_sources ( + source_path TEXT PRIMARY KEY, + canonical_url TEXT, + source_type TEXT, + source_authority INTEGER, + source_hash TEXT, + r2_etag TEXT, + chunk_count INTEGER NOT NULL DEFAULT 0, + approval_status TEXT NOT NULL DEFAULT 'pending', + approved_by TEXT, + approved_at TEXT, + disabled_at TEXT, + verified_at TEXT, + indexed_at TEXT, + created_at TEXT NOT NULL, + updated_at TEXT + ); + CREATE TABLE IF NOT EXISTS retrieval_runs ( id TEXT PRIMARY KEY, answer_run_id TEXT, diff --git a/server/src/db/migrations/0029_knowledge_sources.sql b/server/src/db/migrations/0029_knowledge_sources.sql new file mode 100644 index 00000000..8ded99da --- /dev/null +++ b/server/src/db/migrations/0029_knowledge_sources.sql @@ -0,0 +1,20 @@ +CREATE TABLE knowledge_sources ( + source_path TEXT PRIMARY KEY, + canonical_url TEXT, + source_type TEXT, + source_authority INTEGER, + source_hash TEXT, + r2_etag TEXT, + chunk_count INTEGER NOT NULL DEFAULT 0, + approval_status TEXT NOT NULL DEFAULT 'pending', + approved_by TEXT, + approved_at TEXT, + disabled_at TEXT, + verified_at TEXT, + indexed_at TEXT, + created_at TEXT NOT NULL, + updated_at TEXT +); + +CREATE INDEX idx_knowledge_sources_status ON knowledge_sources(approval_status); +CREATE INDEX idx_knowledge_sources_canonical_url ON knowledge_sources(canonical_url); diff --git a/server/src/db/schema.ts b/server/src/db/schema.ts index b86486d8..fa45936e 100644 --- a/server/src/db/schema.ts +++ b/server/src/db/schema.ts @@ -214,6 +214,27 @@ export const llmUsage = sqliteTable("llm_usage", { export type LlmUsage = typeof llmUsage.$inferSelect; export type NewLlmUsage = typeof llmUsage.$inferInsert; +export const knowledgeSources = sqliteTable("knowledge_sources", { + sourcePath: text("source_path").primaryKey(), + canonicalUrl: text("canonical_url"), + sourceType: text("source_type"), + sourceAuthority: integer("source_authority"), + sourceHash: text("source_hash"), + r2Etag: text("r2_etag"), + chunkCount: integer("chunk_count").notNull().default(0), + approvalStatus: text("approval_status").notNull().default("pending"), // "pending" | "approved" | "rejected" | "disabled" + approvedBy: text("approved_by"), + approvedAt: text("approved_at"), + disabledAt: text("disabled_at"), + verifiedAt: text("verified_at"), + indexedAt: text("indexed_at"), + createdAt: text("created_at").notNull(), + updatedAt: text("updated_at"), +}); + +export type KnowledgeSource = typeof knowledgeSources.$inferSelect; +export type NewKnowledgeSource = typeof knowledgeSources.$inferInsert; + export const retrievalRuns = sqliteTable("retrieval_runs", { id: text("id").primaryKey(), answerRunId: text("answer_run_id"), From 17273f1468d1f56218175fe4332d52a49298d55b Mon Sep 17 00:00:00 2001 From: owk-owk130 Date: Tue, 1 Sep 2026 17:11:25 +0900 Subject: [PATCH 2/6] =?UTF-8?q?refactor(server):=20SHA-256=20=E3=81=AE=20h?= =?UTF-8?q?ex=20=E5=8C=96=E3=82=92=20lib/crypto=20=E3=81=AB=E5=85=B1?= =?UTF-8?q?=E9=80=9A=E5=8C=96=E3=81=99=E3=82=8B?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-Authored-By: Claude Fable 5 Claude-Session: https://claude.ai/code/session_018BjNFvQjDzMwF8HaZHnJjP --- server/src/lib/crypto.test.ts | 14 ++++++++++++++ server/src/lib/crypto.ts | 10 ++++++++++ server/src/services/knowledge/embedding.ts | 13 ++----------- 3 files changed, 26 insertions(+), 11 deletions(-) diff --git a/server/src/lib/crypto.test.ts b/server/src/lib/crypto.test.ts index 446e24ca..312332c2 100644 --- a/server/src/lib/crypto.test.ts +++ b/server/src/lib/crypto.test.ts @@ -7,6 +7,7 @@ import { generateToken, hmacSha1Base64, hmacSha256, + sha256Hex, } from "./crypto"; describe("hmacSha256", () => { @@ -109,3 +110,16 @@ describe("generateToken", () => { expect(tokens.size).toBe(20); }); }); + +describe("sha256Hex", () => { + it("既知の入力に対して SHA-256 の hex を返す", async () => { + expect(await sha256Hex("abc")).toBe( + "ba7816bf8f01cfea414140de5dae2223b00361a396177a9cb410ff61f20015ad", + ); + }); + + it("同じ入力は同じハッシュ、異なる入力は異なるハッシュになる", async () => { + expect(await sha256Hex("foo")).toBe(await sha256Hex("foo")); + expect(await sha256Hex("foo")).not.toBe(await sha256Hex("bar")); + }); +}); diff --git a/server/src/lib/crypto.ts b/server/src/lib/crypto.ts index 224cec84..2651b780 100644 --- a/server/src/lib/crypto.ts +++ b/server/src/lib/crypto.ts @@ -61,6 +61,16 @@ export const hmacSha1Base64 = async (value: string, secret: string) => { return btoa(String.fromCharCode(...new Uint8Array(signature))); }; +export const sha256Hex = async (value: string) => { + const digest = await crypto.subtle.digest( + "SHA-256", + new TextEncoder().encode(value), + ); + return Array.from(new Uint8Array(digest), (byte) => + byte.toString(16).padStart(2, "0"), + ).join(""); +}; + export const generateId = () => { const array = new Uint8Array(16); crypto.getRandomValues(array); diff --git a/server/src/services/knowledge/embedding.ts b/server/src/services/knowledge/embedding.ts index 616cfbee..10d1af74 100644 --- a/server/src/services/knowledge/embedding.ts +++ b/server/src/services/knowledge/embedding.ts @@ -2,6 +2,7 @@ import { createGoogleGenerativeAI } from "@ai-sdk/google"; import { MDocument } from "@mastra/rag"; import { embedMany } from "ai"; import matter from "gray-matter"; +import { sha256Hex } from "~/lib/crypto"; import { GEMINI_EMBEDDING } from "~/lib/llm-models"; import { logger } from "~/lib/logger"; import { recordLlmUsage } from "~/services/analytics/llm-usage"; @@ -25,16 +26,6 @@ type ChunkMetadata = { [key: string]: string | number | boolean | string[] | undefined; }; -const hashText = async (text: string) => { - const digest = await crypto.subtle.digest( - "SHA-256", - new TextEncoder().encode(text), - ); - return [...new Uint8Array(digest)] - .map((byte) => byte.toString(16).padStart(2, "0")) - .join(""); -}; - type VectorData = { id: string; values: number[]; @@ -108,7 +99,7 @@ const chunkDocument = async ( section: chunkMeta?.section as string | undefined, subsection: chunkMeta?.subsection as string | undefined, content: texts[idx], - contentHash: await hashText(texts[idx]), + contentHash: await sha256Hex(texts[idx]), }; }), ); From c96ed05742128da34fa49e0b058e98c7ca164f2f Mon Sep 17 00:00:00 2001 From: owk-owk130 Date: Tue, 1 Sep 2026 17:11:25 +0900 Subject: [PATCH 3/6] =?UTF-8?q?feat(server):=20=E6=83=85=E5=A0=B1=E6=BA=90?= =?UTF-8?q?=E3=83=A1=E3=82=BF=E3=83=87=E3=83=BC=E3=82=BF=E6=8A=BD=E5=87=BA?= =?UTF-8?q?=E3=81=A8=20manifest=20repository=20=E3=82=92=E8=BF=BD=E5=8A=A0?= =?UTF-8?q?=E3=81=99=E3=82=8B?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-Authored-By: Claude Fable 5 Claude-Session: https://claude.ai/code/session_018BjNFvQjDzMwF8HaZHnJjP --- .../knowledge-source-repository.test.ts | 139 ++++++++++++++++++ .../repository/knowledge-source-repository.ts | 80 ++++++++++ .../services/knowledge/source-meta.test.ts | 46 ++++++ server/src/services/knowledge/source-meta.ts | 31 ++++ 4 files changed, 296 insertions(+) create mode 100644 server/src/repository/knowledge-source-repository.test.ts create mode 100644 server/src/repository/knowledge-source-repository.ts create mode 100644 server/src/services/knowledge/source-meta.test.ts create mode 100644 server/src/services/knowledge/source-meta.ts diff --git a/server/src/repository/knowledge-source-repository.test.ts b/server/src/repository/knowledge-source-repository.test.ts new file mode 100644 index 00000000..3b82ca6c --- /dev/null +++ b/server/src/repository/knowledge-source-repository.test.ts @@ -0,0 +1,139 @@ +import { beforeEach, describe, expect, it, vi } from "vitest"; +import { createTestDb, type TestDb } from "~/__tests__/helpers/test-db"; + +const { testDbHolder } = vi.hoisted(() => ({ + testDbHolder: { db: null as TestDb | null }, +})); + +vi.mock("~/db", async (importOriginal) => { + const actual = await importOriginal(); + return { + ...actual, + createDb: () => testDbHolder.db, + }; +}); + +const { knowledgeSourceRepository } = await import( + "./knowledge-source-repository" +); + +const d1 = {} as D1Database; + +const baseSource = { + sourcePath: "bus/index.md", + approvalStatus: "pending", + createdAt: "2026-09-01T00:00:00.000Z", +}; + +describe("knowledgeSourceRepository", () => { + beforeEach(async () => { + testDbHolder.db = await createTestDb(); + }); + + it("insert と findByPath で登録内容を往復できる", async () => { + await knowledgeSourceRepository.insert(d1, { + ...baseSource, + canonicalUrl: "https://example.com/bus", + sourceAuthority: 1, + }); + + const found = await knowledgeSourceRepository.findByPath( + d1, + "bus/index.md", + ); + expect(found).toMatchObject({ + sourcePath: "bus/index.md", + canonicalUrl: "https://example.com/bus", + sourceAuthority: 1, + approvalStatus: "pending", + chunkCount: 0, + }); + }); + + it("存在しない source_path は null を返す", async () => { + expect( + await knowledgeSourceRepository.findByPath(d1, "missing.md"), + ).toBeNull(); + }); + + it("list は source_path 昇順で全件返す", async () => { + await knowledgeSourceRepository.insert(d1, { + ...baseSource, + sourcePath: "b.md", + }); + await knowledgeSourceRepository.insert(d1, { + ...baseSource, + sourcePath: "a.md", + }); + + const rows = await knowledgeSourceRepository.list(d1); + expect(rows.map((r) => r.sourcePath)).toEqual(["a.md", "b.md"]); + }); + + it("update はメタデータと updatedAt だけ更新する", async () => { + await knowledgeSourceRepository.insert(d1, baseSource); + + await knowledgeSourceRepository.update(d1, "bus/index.md", { + sourceHash: "hash-2", + r2Etag: "etag-2", + }); + + const found = await knowledgeSourceRepository.findByPath( + d1, + "bus/index.md", + ); + expect(found).toMatchObject({ + sourceHash: "hash-2", + r2Etag: "etag-2", + approvalStatus: "pending", + }); + expect(found?.updatedAt).not.toBeNull(); + }); + + it("update は承認情報を更新する", async () => { + await knowledgeSourceRepository.insert(d1, baseSource); + + await knowledgeSourceRepository.update(d1, "bus/index.md", { + approvalStatus: "approved", + approvedBy: "admin-1", + approvedAt: "2026-09-02T00:00:00.000Z", + }); + + const found = await knowledgeSourceRepository.findByPath( + d1, + "bus/index.md", + ); + expect(found).toMatchObject({ + approvalStatus: "approved", + approvedBy: "admin-1", + approvedAt: "2026-09-02T00:00:00.000Z", + }); + }); + + it("markIndexed は chunk_count と indexed_at を更新する", async () => { + await knowledgeSourceRepository.insert(d1, baseSource); + + await knowledgeSourceRepository.markIndexed(d1, "bus/index.md", 12); + + const found = await knowledgeSourceRepository.findByPath( + d1, + "bus/index.md", + ); + expect(found?.chunkCount).toBe(12); + expect(found?.indexedAt).not.toBeNull(); + }); + + it("markRemoved は chunk_count を 0 に戻し indexed_at を消す", async () => { + await knowledgeSourceRepository.insert(d1, baseSource); + await knowledgeSourceRepository.markIndexed(d1, "bus/index.md", 12); + + await knowledgeSourceRepository.markRemoved(d1, "bus/index.md"); + + const found = await knowledgeSourceRepository.findByPath( + d1, + "bus/index.md", + ); + expect(found?.chunkCount).toBe(0); + expect(found?.indexedAt).toBeNull(); + }); +}); diff --git a/server/src/repository/knowledge-source-repository.ts b/server/src/repository/knowledge-source-repository.ts new file mode 100644 index 00000000..6b201153 --- /dev/null +++ b/server/src/repository/knowledge-source-repository.ts @@ -0,0 +1,80 @@ +import { asc, eq } from "drizzle-orm"; +import { + createDb, + type KnowledgeSource, + knowledgeSources, + type NewKnowledgeSource, +} from "~/db"; + +export const APPROVAL_STATUSES = [ + "pending", + "approved", + "rejected", + "disabled", +] as const; + +export type ApprovalStatus = (typeof APPROVAL_STATUSES)[number]; + +type UpdateInput = Partial>; + +const update = async ( + d1: D1Database, + sourcePath: string, + patch: UpdateInput, +) => { + const db = createDb(d1); + await db + .update(knowledgeSources) + .set({ ...patch, updatedAt: new Date().toISOString() }) + .where(eq(knowledgeSources.sourcePath, sourcePath)); +}; + +export const knowledgeSourceRepository = { + async findByPath(d1: D1Database, sourcePath: string) { + const db = createDb(d1); + const result = await db + .select() + .from(knowledgeSources) + .where(eq(knowledgeSources.sourcePath, sourcePath)) + .get(); + return result ?? null; + }, + + async list(d1: D1Database) { + const db = createDb(d1); + return await db + .select() + .from(knowledgeSources) + .orderBy(asc(knowledgeSources.sourcePath)) + .all(); + }, + + async insert(d1: D1Database, values: NewKnowledgeSource) { + const db = createDb(d1); + return await db.insert(knowledgeSources).values(values).returning().get(); + }, + + update, + + markIndexed(d1: D1Database, sourcePath: string, chunkCount: number) { + return update(d1, sourcePath, { + chunkCount, + indexedAt: new Date().toISOString(), + }); + }, + + markRemoved(d1: D1Database, sourcePath: string) { + return update(d1, sourcePath, { chunkCount: 0, indexedAt: null }); + }, + + async markAllRemoved(d1: D1Database) { + const db = createDb(d1); + await db.update(knowledgeSources).set({ + chunkCount: 0, + indexedAt: null, + updatedAt: new Date().toISOString(), + }); + }, +}; + +export type { KnowledgeSource }; diff --git a/server/src/services/knowledge/source-meta.test.ts b/server/src/services/knowledge/source-meta.test.ts new file mode 100644 index 00000000..888eb390 --- /dev/null +++ b/server/src/services/knowledge/source-meta.test.ts @@ -0,0 +1,46 @@ +import { describe, expect, it } from "vitest"; +import { extractSourceMeta } from "./source-meta"; + +describe("extractSourceMeta", () => { + it("frontmatter から情報源メタデータを取り出す", () => { + const content = `--- +url: 'https://www.vill.otoineppu.hokkaido.jp/' +source_type: curated +source_authority: 2 +verified_at: '2026-09-01' +--- +# 本文 +`; + expect(extractSourceMeta(content)).toEqual({ + canonicalUrl: "https://www.vill.otoineppu.hokkaido.jp/", + sourceType: "curated", + sourceAuthority: 2, + verifiedAt: "2026-09-01", + }); + }); + + it("YAML が日付型に解釈した verified_at も日付文字列にする", () => { + const content = `--- +verified_at: 2026-09-01 +--- +本文`; + expect(extractSourceMeta(content).verifiedAt).toBe("2026-09-01"); + }); + + it("キーが無ければ undefined を返す", () => { + expect(extractSourceMeta("# 見出しのみ")).toEqual({ + canonicalUrl: undefined, + sourceType: undefined, + sourceAuthority: undefined, + verifiedAt: undefined, + }); + }); + + it("数値に解釈できない source_authority は undefined にする", () => { + const content = `--- +source_authority: high +--- +本文`; + expect(extractSourceMeta(content).sourceAuthority).toBeUndefined(); + }); +}); diff --git a/server/src/services/knowledge/source-meta.ts b/server/src/services/knowledge/source-meta.ts new file mode 100644 index 00000000..6e8751ba --- /dev/null +++ b/server/src/services/knowledge/source-meta.ts @@ -0,0 +1,31 @@ +import matter from "gray-matter"; +import { sha256Hex } from "~/lib/crypto"; + +const asDateString = (value: unknown) => { + if (typeof value === "string") return value; + if (value instanceof Date) return value.toISOString().slice(0, 10); + return undefined; +}; + +export const extractSourceMeta = (content: string) => { + const { data } = matter(content); + const authority = Number(data.source_authority); + return { + canonicalUrl: typeof data.url === "string" ? data.url : undefined, + sourceType: + typeof data.source_type === "string" ? data.source_type : undefined, + sourceAuthority: Number.isFinite(authority) ? authority : undefined, + verifiedAt: asDateString(data.verified_at), + }; +}; + +export type SourceMeta = ReturnType; + +export const buildSourceRecord = async ( + sourcePath: string, + content: string, +) => ({ + sourcePath, + ...extractSourceMeta(content), + sourceHash: await sha256Hex(content), +}); From 4e7763b06fc2bd6feae99302425b56c26272cc34 Mon Sep 17 00:00:00 2001 From: owk-owk130 Date: Tue, 1 Sep 2026 17:11:25 +0900 Subject: [PATCH 4/6] =?UTF-8?q?feat(server):=20=E6=89=BF=E8=AA=8D=E3=82=B2?= =?UTF-8?q?=E3=83=BC=E3=83=88=E4=BB=98=E3=81=8D=E3=81=AE=E3=83=8A=E3=83=AC?= =?UTF-8?q?=E3=83=83=E3=82=B8=20index=20=E5=87=A6=E7=90=86=E3=82=92?= =?UTF-8?q?=E8=BF=BD=E5=8A=A0=E3=81=99=E3=82=8B?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-Authored-By: Claude Fable 5 Claude-Session: https://claude.ai/code/session_018BjNFvQjDzMwF8HaZHnJjP --- .../src/services/knowledge/indexing.test.ts | 286 ++++++++++++++++++ server/src/services/knowledge/indexing.ts | 117 +++++++ 2 files changed, 403 insertions(+) create mode 100644 server/src/services/knowledge/indexing.test.ts create mode 100644 server/src/services/knowledge/indexing.ts diff --git a/server/src/services/knowledge/indexing.test.ts b/server/src/services/knowledge/indexing.test.ts new file mode 100644 index 00000000..18134b87 --- /dev/null +++ b/server/src/services/knowledge/indexing.test.ts @@ -0,0 +1,286 @@ +import { beforeEach, describe, expect, it, vi } from "vitest"; +import { createTestDb, type TestDb } from "~/__tests__/helpers/test-db"; + +const { testDbHolder } = vi.hoisted(() => ({ + testDbHolder: { db: null as TestDb | null }, +})); + +vi.mock("~/db", async (importOriginal) => { + const actual = await importOriginal(); + return { + ...actual, + createDb: () => testDbHolder.db, + }; +}); + +vi.mock("~/lib/logger", () => ({ + logger: { info: vi.fn(), warn: vi.fn(), error: vi.fn() }, +})); + +vi.mock("./embedding", () => ({ + deleteKnowledgeBySource: vi.fn(async () => ({ deleted: 0 })), + processKnowledgeFile: vi.fn(async () => ({ chunks: 3 })), +})); + +const { deleteKnowledgeBySource, processKnowledgeFile } = await import( + "./embedding" +); +const { sha256Hex } = await import("~/lib/crypto"); +const { knowledgeSourceRepository } = await import( + "~/repository/knowledge-source-repository" +); +const { indexKnowledgeSource, removeKnowledgeSource } = await import( + "./indexing" +); + +const d1 = {} as D1Database; +const vectorize = {} as VectorizeIndex; +const deps = { d1, vectorize, apiKey: "key" }; + +const content = `--- +url: 'https://example.com/bus' +--- +# バス +本文`; + +describe("indexKnowledgeSource", () => { + beforeEach(async () => { + testDbHolder.db = await createTestDb(); + vi.mocked(deleteKnowledgeBySource).mockClear(); + vi.mocked(processKnowledgeFile).mockClear(); + }); + + it("未登録の情報源は pending で登録し index しない", async () => { + const result = await indexKnowledgeSource("bus/index.md", content, deps); + + expect(result).toEqual({ indexed: false, status: "pending", chunks: 0 }); + expect(processKnowledgeFile).not.toHaveBeenCalled(); + expect(deleteKnowledgeBySource).not.toHaveBeenCalled(); + + const row = await knowledgeSourceRepository.findByPath(d1, "bus/index.md"); + expect(row).toMatchObject({ + approvalStatus: "pending", + canonicalUrl: "https://example.com/bus", + }); + expect(row?.sourceHash).toMatch(/^[0-9a-f]{64}$/); + }); + + it("approved の情報源は index して chunk_count を記録する", async () => { + await knowledgeSourceRepository.insert(d1, { + sourcePath: "bus/index.md", + approvalStatus: "approved", + createdAt: "2026-09-01T00:00:00.000Z", + }); + + const result = await indexKnowledgeSource("bus/index.md", content, deps, { + r2Etag: "etag-1", + }); + + expect(result).toMatchObject({ + indexed: true, + status: "approved", + chunks: 3, + }); + expect(deleteKnowledgeBySource).toHaveBeenCalledWith( + vectorize, + "bus/index.md", + ); + + const row = await knowledgeSourceRepository.findByPath(d1, "bus/index.md"); + expect(row).toMatchObject({ chunkCount: 3, r2Etag: "etag-1" }); + expect(row?.indexedAt).not.toBeNull(); + }); + + it("rejected / disabled は index せずメタデータだけ更新する", async () => { + await knowledgeSourceRepository.insert(d1, { + sourcePath: "bus/index.md", + approvalStatus: "rejected", + createdAt: "2026-09-01T00:00:00.000Z", + }); + + const result = await indexKnowledgeSource("bus/index.md", content, deps); + + expect(result).toEqual({ indexed: false, status: "rejected", chunks: 0 }); + expect(processKnowledgeFile).not.toHaveBeenCalled(); + + const row = await knowledgeSourceRepository.findByPath(d1, "bus/index.md"); + expect(row?.canonicalUrl).toBe("https://example.com/bus"); + }); + + it("approveAs 指定で未登録の情報源を approved として登録し index する", async () => { + const result = await indexKnowledgeSource("bus/index.md", content, deps, { + approveAs: "admin-1", + }); + + expect(result).toMatchObject({ indexed: true, chunks: 3 }); + + const row = await knowledgeSourceRepository.findByPath(d1, "bus/index.md"); + expect(row).toMatchObject({ + approvalStatus: "approved", + approvedBy: "admin-1", + }); + expect(row?.approvedAt).not.toBeNull(); + }); + + it("approveAs 指定で pending を approved に昇格して index する", async () => { + await knowledgeSourceRepository.insert(d1, { + sourcePath: "bus/index.md", + approvalStatus: "pending", + createdAt: "2026-09-01T00:00:00.000Z", + }); + + const result = await indexKnowledgeSource("bus/index.md", content, deps, { + approveAs: "admin-1", + }); + + expect(result).toMatchObject({ indexed: true, chunks: 3 }); + + const row = await knowledgeSourceRepository.findByPath(d1, "bus/index.md"); + expect(row).toMatchObject({ + approvalStatus: "approved", + approvedBy: "admin-1", + }); + }); + + it("approveAs 指定でも rejected / disabled は昇格しない", async () => { + await knowledgeSourceRepository.insert(d1, { + sourcePath: "bus/index.md", + approvalStatus: "disabled", + createdAt: "2026-09-01T00:00:00.000Z", + }); + + const result = await indexKnowledgeSource("bus/index.md", content, deps, { + approveAs: "admin-1", + }); + + expect(result).toEqual({ indexed: false, status: "disabled", chunks: 0 }); + expect(processKnowledgeFile).not.toHaveBeenCalled(); + }); + + it("skipUnchanged 指定で内容不変の approved は再 index しない", async () => { + await knowledgeSourceRepository.insert(d1, { + sourcePath: "bus/index.md", + approvalStatus: "approved", + sourceHash: await sha256Hex(content), + chunkCount: 8, + indexedAt: "2026-09-01T00:00:00.000Z", + createdAt: "2026-09-01T00:00:00.000Z", + }); + + const result = await indexKnowledgeSource("bus/index.md", content, deps, { + skipUnchanged: true, + }); + + expect(result).toEqual({ indexed: true, status: "approved", chunks: 8 }); + expect(processKnowledgeFile).not.toHaveBeenCalled(); + expect(deleteKnowledgeBySource).not.toHaveBeenCalled(); + }); + + it("skipUnchanged 指定でも内容が変わっていれば再 index する", async () => { + await knowledgeSourceRepository.insert(d1, { + sourcePath: "bus/index.md", + approvalStatus: "approved", + sourceHash: "old-hash", + indexedAt: "2026-09-01T00:00:00.000Z", + createdAt: "2026-09-01T00:00:00.000Z", + }); + + const result = await indexKnowledgeSource("bus/index.md", content, deps, { + skipUnchanged: true, + }); + + expect(result).toMatchObject({ indexed: true, chunks: 3 }); + expect(processKnowledgeFile).toHaveBeenCalled(); + }); + + it("内容不変なら再 index はしてもメタデータは書き換えない", async () => { + await knowledgeSourceRepository.insert(d1, { + sourcePath: "bus/index.md", + approvalStatus: "approved", + sourceHash: await sha256Hex(content), + canonicalUrl: "https://example.com/manual", + createdAt: "2026-09-01T00:00:00.000Z", + }); + + await indexKnowledgeSource("bus/index.md", content, deps); + + const row = await knowledgeSourceRepository.findByPath(d1, "bus/index.md"); + expect(row?.canonicalUrl).toBe("https://example.com/manual"); + expect(row?.chunkCount).toBe(3); + }); + + it("index 済みなのに未承認の情報源はベクトルを削除して整合させる", async () => { + await knowledgeSourceRepository.insert(d1, { + sourcePath: "bus/index.md", + approvalStatus: "pending", + chunkCount: 8, + indexedAt: "2026-09-01T00:00:00.000Z", + createdAt: "2026-09-01T00:00:00.000Z", + }); + + const result = await indexKnowledgeSource("bus/index.md", content, deps); + + expect(result).toEqual({ indexed: false, status: "pending", chunks: 0 }); + expect(deleteKnowledgeBySource).toHaveBeenCalledWith( + vectorize, + "bus/index.md", + ); + expect(processKnowledgeFile).not.toHaveBeenCalled(); + + const row = await knowledgeSourceRepository.findByPath(d1, "bus/index.md"); + expect(row?.chunkCount).toBe(0); + expect(row?.indexedAt).toBeNull(); + }); + + it("index が失敗したら chunk_count を更新せず error を返す", async () => { + vi.mocked(processKnowledgeFile).mockResolvedValueOnce({ + chunks: 0, + error: "embedding failed", + }); + await knowledgeSourceRepository.insert(d1, { + sourcePath: "bus/index.md", + approvalStatus: "approved", + createdAt: "2026-09-01T00:00:00.000Z", + }); + + const result = await indexKnowledgeSource("bus/index.md", content, deps); + + expect(result).toMatchObject({ indexed: true, error: "embedding failed" }); + + const row = await knowledgeSourceRepository.findByPath(d1, "bus/index.md"); + expect(row?.indexedAt).toBeNull(); + }); +}); + +describe("removeKnowledgeSource", () => { + beforeEach(async () => { + testDbHolder.db = await createTestDb(); + vi.mocked(deleteKnowledgeBySource).mockClear(); + }); + + it("ベクトルを削除し manifest の chunk_count を 0 に戻す", async () => { + await knowledgeSourceRepository.insert(d1, { + sourcePath: "bus/index.md", + approvalStatus: "approved", + createdAt: "2026-09-01T00:00:00.000Z", + }); + await knowledgeSourceRepository.markIndexed(d1, "bus/index.md", 5); + + await removeKnowledgeSource("bus/index.md", { d1, vectorize }); + + expect(deleteKnowledgeBySource).toHaveBeenCalledWith( + vectorize, + "bus/index.md", + ); + const row = await knowledgeSourceRepository.findByPath(d1, "bus/index.md"); + expect(row?.chunkCount).toBe(0); + }); + + it("manifest 行が無くてもベクトル削除は行う", async () => { + await removeKnowledgeSource("unknown.md", { d1, vectorize }); + expect(deleteKnowledgeBySource).toHaveBeenCalledWith( + vectorize, + "unknown.md", + ); + }); +}); diff --git a/server/src/services/knowledge/indexing.ts b/server/src/services/knowledge/indexing.ts new file mode 100644 index 00000000..35fb194a --- /dev/null +++ b/server/src/services/knowledge/indexing.ts @@ -0,0 +1,117 @@ +import { logger } from "~/lib/logger"; +import { knowledgeSourceRepository } from "~/repository/knowledge-source-repository"; +import { deleteKnowledgeBySource, processKnowledgeFile } from "./embedding"; +import { buildSourceRecord } from "./source-meta"; + +type IndexDeps = { + d1: D1Database; + vectorize: VectorizeIndex; + apiKey: string; +}; + +type IndexOptions = { + r2Etag?: string; + approveAs?: string; + skipUnchanged?: boolean; +}; + +export type IndexResult = { + indexed: boolean; + status: string; + chunks: number; + error?: string; +}; + +export const indexKnowledgeSource = async ( + key: string, + content: string, + deps: IndexDeps, + options: IndexOptions = {}, +): Promise => { + const record = await buildSourceRecord(key, content); + const now = new Date().toISOString(); + const existing = await knowledgeSourceRepository.findByPath(deps.d1, key); + + const promoted = + options.approveAs !== undefined && + (!existing || existing.approvalStatus === "pending"); + const status = promoted + ? "approved" + : (existing?.approvalStatus ?? "pending"); + + if ( + options.skipUnchanged && + existing?.approvalStatus === "approved" && + existing.indexedAt && + existing.sourceHash === record.sourceHash + ) { + logger.info(`[Knowledge] skip indexing (unchanged): ${key}`); + return { indexed: true, status, chunks: existing.chunkCount }; + } + + if (!existing) { + await knowledgeSourceRepository.insert(deps.d1, { + ...record, + r2Etag: options.r2Etag, + approvalStatus: status, + approvedBy: promoted ? options.approveAs : undefined, + approvedAt: promoted ? now : undefined, + createdAt: now, + }); + } else { + if (existing.sourceHash !== record.sourceHash) { + const { sourcePath: _, ...meta } = record; + await knowledgeSourceRepository.update(deps.d1, key, { + ...meta, + ...(options.r2Etag && { r2Etag: options.r2Etag }), + }); + } + if (promoted && existing.approvalStatus === "pending") { + await knowledgeSourceRepository.update(deps.d1, key, { + approvalStatus: "approved", + approvedBy: options.approveAs, + approvedAt: now, + }); + } + } + + if (status !== "approved") { + if (existing?.indexedAt) { + await deleteKnowledgeBySource(deps.vectorize, key); + await knowledgeSourceRepository.markRemoved(deps.d1, key); + } + logger.info(`[Knowledge] skip indexing (${status}): ${key}`); + return { indexed: false, status, chunks: 0 }; + } + + await deleteKnowledgeBySource(deps.vectorize, key); + const result = await processKnowledgeFile( + key, + content, + deps.vectorize, + deps.apiKey, + deps.d1, + ); + + if (!result.error) { + await knowledgeSourceRepository.markIndexed(deps.d1, key, result.chunks); + } + + return { + indexed: true, + status, + chunks: result.chunks, + ...(result.error !== undefined && { error: result.error }), + }; +}; + +export const removeKnowledgeSource = async ( + key: string, + deps: Pick, +) => { + const result = await deleteKnowledgeBySource(deps.vectorize, key); + if (await knowledgeSourceRepository.findByPath(deps.d1, key)) { + await knowledgeSourceRepository.markRemoved(deps.d1, key); + } + return result; +}; From ef13f11621df8021a5a0fd579c0e52a4f8b2870b Mon Sep 17 00:00:00 2001 From: owk-owk130 Date: Tue, 1 Sep 2026 17:11:26 +0900 Subject: [PATCH 5/6] =?UTF-8?q?refactor(server):=20=E5=85=A8=20index=20?= =?UTF-8?q?=E7=B5=8C=E8=B7=AF=E3=82=92=E6=89=BF=E8=AA=8D=E3=82=B2=E3=83=BC?= =?UTF-8?q?=E3=83=88=E7=B5=8C=E7=94=B1=E3=81=AB=E3=81=99=E3=82=8B?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-Authored-By: Claude Fable 5 Claude-Session: https://claude.ai/code/session_018BjNFvQjDzMwF8HaZHnJjP --- server/src/handlers/r2-event-handler.test.ts | 67 +++++++---- server/src/handlers/r2-event-handler.ts | 42 ++++--- server/src/routes/admin/knowledge/convert.ts | 4 + server/src/routes/admin/knowledge/files.ts | 8 +- .../src/routes/admin/knowledge/index.test.ts | 19 +++- server/src/routes/admin/knowledge/sync.ts | 2 + server/src/services/knowledge/files.test.ts | 21 ++-- server/src/services/knowledge/files.ts | 5 +- server/src/services/knowledge/sync.test.ts | 105 +++++++++++------- server/src/services/knowledge/sync.ts | 45 ++++---- server/src/services/knowledge/upload.test.ts | 39 ++++++- server/src/services/knowledge/upload.ts | 6 +- 12 files changed, 255 insertions(+), 108 deletions(-) diff --git a/server/src/handlers/r2-event-handler.test.ts b/server/src/handlers/r2-event-handler.test.ts index 31593c1e..ea596c45 100644 --- a/server/src/handlers/r2-event-handler.test.ts +++ b/server/src/handlers/r2-event-handler.test.ts @@ -1,12 +1,12 @@ import { beforeEach, describe, expect, it, vi } from "vitest"; -vi.mock("~/services/knowledge/embedding", () => ({ - deleteKnowledgeBySource: vi.fn(), - processKnowledgeFile: vi.fn(), +vi.mock("~/services/knowledge/indexing", () => ({ + indexKnowledgeSource: vi.fn(), + removeKnowledgeSource: vi.fn(), })); -const { deleteKnowledgeBySource, processKnowledgeFile } = await import( - "~/services/knowledge/embedding" +const { indexKnowledgeSource, removeKnowledgeSource } = await import( + "~/services/knowledge/indexing" ); const { handleR2Event } = await import("./r2-event-handler"); @@ -15,6 +15,7 @@ const r2Bucket = { }; const env = { + DB: {} as D1Database, KNOWLEDGE_BUCKET: r2Bucket, VECTORIZE: {} as VectorizeIndex, GOOGLE_GENERATIVE_AI_API_KEY: "key", @@ -56,8 +57,12 @@ describe("handleR2Event", () => { beforeEach(() => { vi.clearAllMocks(); r2Bucket.get.mockResolvedValue({ text: vi.fn().mockResolvedValue("md") }); - vi.mocked(deleteKnowledgeBySource).mockResolvedValue({ deleted: 3 }); - vi.mocked(processKnowledgeFile).mockResolvedValue({ chunks: 5 }); + vi.mocked(indexKnowledgeSource).mockResolvedValue({ + indexed: true, + status: "approved", + chunks: 5, + }); + vi.mocked(removeKnowledgeSource).mockResolvedValue({ deleted: 3 }); }); it(".md 以外は ack して何もしない", async () => { @@ -66,34 +71,56 @@ describe("handleR2Event", () => { await handleR2Event(buildBatch([m]), env); expect(m.ack).toHaveBeenCalled(); - expect(processKnowledgeFile).not.toHaveBeenCalled(); + expect(indexKnowledgeSource).not.toHaveBeenCalled(); }); it.each(["PutObject", "CompleteMultipartUpload", "CopyObject"] as const)( - "%s は delete + processKnowledgeFile を順に呼ぶ", + "%s は indexKnowledgeSource に eTag 付きで委譲する", async (action) => { const m = buildMessage(action, "doc.md"); await handleR2Event(buildBatch([m]), env); - expect(deleteKnowledgeBySource).toHaveBeenCalledWith( - env.VECTORIZE, + expect(indexKnowledgeSource).toHaveBeenCalledWith( "doc.md", + "md", + { + d1: env.DB, + vectorize: env.VECTORIZE, + apiKey: "key", + }, + { r2Etag: "etag", skipUnchanged: true }, ); - expect(processKnowledgeFile).toHaveBeenCalled(); expect(m.ack).toHaveBeenCalled(); }, ); + it("未承認で index が skip されても ack する", async () => { + vi.mocked(indexKnowledgeSource).mockResolvedValue({ + indexed: false, + status: "pending", + chunks: 0, + }); + const m = buildMessage("PutObject", "doc.md"); + + await handleR2Event(buildBatch([m]), env); + + expect(m.ack).toHaveBeenCalled(); + expect(m.retry).not.toHaveBeenCalled(); + }); + it.each(["DeleteObject", "LifecycleDeletion"] as const)( - "%s は deleteKnowledgeBySource のみ", + "%s は removeKnowledgeSource のみ", async (action) => { const m = buildMessage(action, "doc.md"); await handleR2Event(buildBatch([m]), env); - expect(deleteKnowledgeBySource).toHaveBeenCalled(); - expect(processKnowledgeFile).not.toHaveBeenCalled(); + expect(removeKnowledgeSource).toHaveBeenCalledWith("doc.md", { + d1: env.DB, + vectorize: env.VECTORIZE, + }); + expect(indexKnowledgeSource).not.toHaveBeenCalled(); expect(m.ack).toHaveBeenCalled(); }, ); @@ -107,10 +134,12 @@ describe("handleR2Event", () => { expect(m.retry).toHaveBeenCalled(); }); - it("processKnowledgeFile がエラーを返したら retry", async () => { - vi.mocked(processKnowledgeFile).mockResolvedValue({ - error: "embed failed", + it("index がエラーを返したら retry", async () => { + vi.mocked(indexKnowledgeSource).mockResolvedValue({ + indexed: true, + status: "approved", chunks: 0, + error: "embed failed", }); const m = buildMessage("PutObject", "doc.md"); @@ -120,7 +149,7 @@ describe("handleR2Event", () => { }); it("例外が起きたら retry", async () => { - vi.mocked(processKnowledgeFile).mockRejectedValue(new Error("boom")); + vi.mocked(indexKnowledgeSource).mockRejectedValue(new Error("boom")); const m = buildMessage("PutObject", "doc.md"); await handleR2Event(buildBatch([m]), env); diff --git a/server/src/handlers/r2-event-handler.ts b/server/src/handlers/r2-event-handler.ts index 6a536534..ac680d82 100644 --- a/server/src/handlers/r2-event-handler.ts +++ b/server/src/handlers/r2-event-handler.ts @@ -1,9 +1,9 @@ import * as Sentry from "@sentry/cloudflare"; import { logger } from "~/lib/logger"; import { - deleteKnowledgeBySource, - processKnowledgeFile, -} from "~/services/knowledge/embedding"; + indexKnowledgeSource, + removeKnowledgeSource, +} from "~/services/knowledge/indexing"; type R2EventType = | "PutObject" @@ -30,8 +30,14 @@ const isMarkdownFile = (key: string) => key.endsWith(".md"); const handleObjectCreate = async ( key: string, + eTag: string, env: CloudflareBindings, -): Promise<{ success: boolean; chunks?: number; error?: string }> => { +): Promise<{ + success: boolean; + chunks?: number; + skipped?: string; + error?: string; +}> => { const file = await env.KNOWLEDGE_BUCKET.get(key); if (!file) { return { success: false, error: `File not found: ${key}` }; @@ -39,16 +45,21 @@ const handleObjectCreate = async ( const content = await file.text(); - // 既存データを削除してから再登録 - await deleteKnowledgeBySource(env.VECTORIZE, key); - - const result = await processKnowledgeFile( + const result = await indexKnowledgeSource( key, content, - env.VECTORIZE, - env.GOOGLE_GENERATIVE_AI_API_KEY, + { + d1: env.DB, + vectorize: env.VECTORIZE, + apiKey: env.GOOGLE_GENERATIVE_AI_API_KEY, + }, + { r2Etag: eTag, skipUnchanged: true }, ); + if (!result.indexed) { + return { success: true, skipped: result.status }; + } + if (result.error) { return { success: false, error: result.error }; } @@ -61,7 +72,10 @@ const handleObjectDelete = async ( env: CloudflareBindings, ): Promise<{ success: boolean; deleted?: number; error?: string }> => { try { - const result = await deleteKnowledgeBySource(env.VECTORIZE, key); + const result = await removeKnowledgeSource(key, { + d1: env.DB, + vectorize: env.VECTORIZE, + }); return { success: true, deleted: result.deleted }; } catch (error) { return { @@ -94,8 +108,10 @@ export const handleR2Event = async ( case "PutObject": case "CompleteMultipartUpload": case "CopyObject": { - const result = await handleObjectCreate(key, env); - if (result.success) { + const result = await handleObjectCreate(key, object.eTag, env); + if (result.success && result.skipped) { + logger.info(`Skipped ${key}: approval status is ${result.skipped}`); + } else if (result.success) { logger.info(`Synced ${key}: ${result.chunks} chunks`); } else { logger.error(`Failed to sync ${key}`, result.error); diff --git a/server/src/routes/admin/knowledge/convert.ts b/server/src/routes/admin/knowledge/convert.ts index 16eb435a..396a8a81 100644 --- a/server/src/routes/admin/knowledge/convert.ts +++ b/server/src/routes/admin/knowledge/convert.ts @@ -3,6 +3,7 @@ import { HTTPException } from "hono/http-exception"; import { errorResponse } from "~/lib/openapi-errors"; import type { PrincipalVariables } from "~/lib/principal"; +import { requireAdminUser } from "~/lib/principal"; import { convertAndUpload, reconvertFromOriginal, @@ -73,6 +74,7 @@ knowledgeConvertRoutes.openapi(uploadFileRoute, async (c) => { vectorize: c.env.VECTORIZE, apiKey, d1: c.env.DB, + approveAs: requireAdminUser(c.get("principal")).id, }); return c.json( @@ -148,6 +150,7 @@ knowledgeConvertRoutes.openapi(convertFileRoute, async (c) => { vectorize: c.env.VECTORIZE, apiKey, d1: c.env.DB, + approveAs: requireAdminUser(c.get("principal")).id, }); return c.json( @@ -219,6 +222,7 @@ knowledgeConvertRoutes.openapi(reconvertFileRoute, async (c) => { vectorize: c.env.VECTORIZE, apiKey, d1: c.env.DB, + approveAs: requireAdminUser(c.get("principal")).id, }); return c.json( diff --git a/server/src/routes/admin/knowledge/files.ts b/server/src/routes/admin/knowledge/files.ts index 7c4a1aca..c0587798 100644 --- a/server/src/routes/admin/knowledge/files.ts +++ b/server/src/routes/admin/knowledge/files.ts @@ -3,6 +3,7 @@ import { HTTPException } from "hono/http-exception"; import { errorResponse } from "~/lib/openapi-errors"; import type { PrincipalVariables } from "~/lib/principal"; +import { requireAdminUser } from "~/lib/principal"; import { deleteFile, getFile, @@ -125,11 +126,14 @@ knowledgeFilesRoutes.openapi(saveFileRoute, async (c) => { vectorize: c.env.VECTORIZE, apiKey, d1: c.env.DB, + approveAs: requireAdminUser(c.get("principal")).id, }); return c.json( { - message: `ファイルを保存し、${result.chunks}チャンクを同期しました`, + message: result.indexed + ? `ファイルを保存し、${result.chunks}チャンクを同期しました` + : `ファイルを保存しました。情報源が ${result.status} のため検索には反映されません`, chunks: result.chunks, }, 200, @@ -159,7 +163,7 @@ knowledgeFilesRoutes.openapi(deleteFileRoute, async (c) => { const { key } = c.req.valid("param"); validateFileKey(key); - await deleteFile(c.env.KNOWLEDGE_BUCKET, c.env.VECTORIZE, key); + await deleteFile(c.env.KNOWLEDGE_BUCKET, c.env.VECTORIZE, c.env.DB, key); const baseName = key.replace(/\.md$/, ""); return c.json({ message: `${baseName} を完全に削除しました` }, 200); }); diff --git a/server/src/routes/admin/knowledge/index.test.ts b/server/src/routes/admin/knowledge/index.test.ts index d55a0487..c9da54c6 100644 --- a/server/src/routes/admin/knowledge/index.test.ts +++ b/server/src/routes/admin/knowledge/index.test.ts @@ -17,6 +17,17 @@ vi.mock("~/services/knowledge", () => ({ reconvertFromOriginal: vi.fn(), })); +vi.mock("~/repository/knowledge-source-repository", () => ({ + APPROVAL_STATUSES: ["pending", "approved", "rejected", "disabled"], + knowledgeSourceRepository: { + findByPath: vi.fn(), + list: vi.fn(), + insert: vi.fn(), + update: vi.fn(), + markAllRemoved: vi.fn(), + }, +})); + vi.mock("~/repository/admin-session-repository", () => ({ adminSessionRepository: { findValid: vi.fn(), @@ -241,6 +252,7 @@ describe("knowledge routes 統合テスト", () => { { file: "b.md", chunks: 5 }, ], editedCount: 0, + skippedCount: 0, }); const res = await app.request( @@ -330,7 +342,12 @@ describe("knowledge routes 統合テスト", () => { }); it("正常系: bucket.put → syncFile → 200", async () => { - vi.mocked(knowledgeService.syncFile).mockResolvedValue({ chunks: 4 }); + vi.mocked(knowledgeService.syncFile).mockResolvedValue({ + indexed: true, + status: "approved", + chunks: 4, + error: undefined, + }); const res = await app.request( authedRequest("/files/doc.md", jsonBody({ content: "# c" })), diff --git a/server/src/routes/admin/knowledge/sync.ts b/server/src/routes/admin/knowledge/sync.ts index c8aeb05b..48f2e01f 100644 --- a/server/src/routes/admin/knowledge/sync.ts +++ b/server/src/routes/admin/knowledge/sync.ts @@ -2,6 +2,7 @@ import { createRoute, OpenAPIHono, z } from "@hono/zod-openapi"; import { errorResponse } from "~/lib/openapi-errors"; import type { PrincipalVariables } from "~/lib/principal"; +import { knowledgeSourceRepository } from "~/repository/knowledge-source-repository"; import { deleteAllKnowledge, syncAll } from "~/services/knowledge"; import { requireApiKey, SuccessResponseSchema } from "./schemas"; @@ -29,6 +30,7 @@ const deleteAllRoute = createRoute({ knowledgeSyncRoutes.openapi(deleteAllRoute, async (c) => { const result = await deleteAllKnowledge(c.env.VECTORIZE); + await knowledgeSourceRepository.markAllRemoved(c.env.DB); return c.json( { message: `${result.deleted}件のベクトルを削除しました`, diff --git a/server/src/services/knowledge/files.test.ts b/server/src/services/knowledge/files.test.ts index 6dd6c6c5..cb6d07d4 100644 --- a/server/src/services/knowledge/files.test.ts +++ b/server/src/services/knowledge/files.test.ts @@ -1,14 +1,14 @@ import { beforeEach, describe, expect, it, vi } from "vitest"; -vi.mock("./embedding", () => ({ - deleteKnowledgeBySource: vi.fn(), +vi.mock("./indexing", () => ({ + removeKnowledgeSource: vi.fn(), })); vi.mock("~/lib/logger", () => ({ logger: { info: vi.fn(), warn: vi.fn(), error: vi.fn() }, })); -const { deleteKnowledgeBySource } = await import("./embedding"); +const { removeKnowledgeSource } = await import("./indexing"); const { deleteFile, getFile, getOriginalFile, listFiles, listUnifiedFiles } = await import("./files"); @@ -59,7 +59,7 @@ const buildBucket = ( }; beforeEach(() => { - vi.mocked(deleteKnowledgeBySource).mockReset(); + vi.mocked(removeKnowledgeSource).mockReset(); }); describe("listFiles", () => { @@ -287,6 +287,8 @@ describe("getOriginalFile", () => { }); }); +const d1 = {} as D1Database; + describe("deleteFile", () => { it("Markdown + originals + Vectorize を削除", async () => { const bucket = buildBucket([ @@ -298,18 +300,21 @@ describe("deleteFile", () => { ]); const vectorize = {} as VectorizeIndex; - await deleteFile(bucket, vectorize, "doc.md"); + await deleteFile(bucket, vectorize, d1, "doc.md"); expect(bucket.delete).toHaveBeenCalledWith("doc.md"); expect(bucket.delete).toHaveBeenCalledWith("originals/doc.pdf"); - expect(deleteKnowledgeBySource).toHaveBeenCalledWith(vectorize, "doc.md"); + expect(removeKnowledgeSource).toHaveBeenCalledWith("doc.md", { + d1, + vectorize, + }); }); it("key に拡張子無しを渡しても .md を付けて削除", async () => { const bucket = buildBucket([]); const vectorize = {} as VectorizeIndex; - await deleteFile(bucket, vectorize, "doc"); + await deleteFile(bucket, vectorize, d1, "doc"); expect(bucket.delete).toHaveBeenCalledWith("doc.md"); }); @@ -324,7 +329,7 @@ describe("deleteFile", () => { ]); const vectorize = {} as VectorizeIndex; - await deleteFile(bucket, vectorize, "doc.md"); + await deleteFile(bucket, vectorize, d1, "doc.md"); // doc.md は削除される expect(bucket.delete).toHaveBeenCalledWith("doc.md"); diff --git a/server/src/services/knowledge/files.ts b/server/src/services/knowledge/files.ts index 636c4d02..456c30d2 100644 --- a/server/src/services/knowledge/files.ts +++ b/server/src/services/knowledge/files.ts @@ -1,5 +1,5 @@ import { logger } from "~/lib/logger"; -import { deleteKnowledgeBySource } from "./embedding"; +import { removeKnowledgeSource } from "./indexing"; import { buildOriginalsMap, EDIT_THRESHOLD_MS, extractBaseName } from "./utils"; export type FileInfo = { @@ -188,6 +188,7 @@ export const getOriginalFile = async ( export const deleteFile = async ( bucket: R2Bucket, vectorize: VectorizeIndex, + d1: D1Database, key: string, ): Promise => { const baseName = key.replace(/\.md$/, ""); @@ -208,6 +209,6 @@ export const deleteFile = async ( } // 3. Vectorize から削除 - await deleteKnowledgeBySource(vectorize, mdKey); + await removeKnowledgeSource(mdKey, { d1, vectorize }); logger.info(`[Delete] Deleted ${mdKey} from Vectorize`); }; diff --git a/server/src/services/knowledge/sync.test.ts b/server/src/services/knowledge/sync.test.ts index edbdd504..8ddbf8dc 100644 --- a/server/src/services/knowledge/sync.test.ts +++ b/server/src/services/knowledge/sync.test.ts @@ -1,17 +1,14 @@ import { beforeEach, describe, expect, it, vi } from "vitest"; -vi.mock("./embedding", () => ({ - deleteKnowledgeBySource: vi.fn(async () => ({ deleted: 0 })), - processKnowledgeFile: vi.fn(), +vi.mock("./indexing", () => ({ + indexKnowledgeSource: vi.fn(), })); vi.mock("~/lib/logger", () => ({ logger: { info: vi.fn(), warn: vi.fn(), error: vi.fn() }, })); -const { deleteKnowledgeBySource, processKnowledgeFile } = await import( - "./embedding" -); +const { indexKnowledgeSource } = await import("./indexing"); const { syncAll, syncFile } = await import("./sync"); const buildR2Object = ( @@ -20,6 +17,7 @@ const buildR2Object = ( ): R2Object => ({ key, + etag: `etag-${key}`, uploaded: uploadedAt instanceof Date ? uploadedAt : new Date(uploadedAt), }) as unknown as R2Object; @@ -40,23 +38,24 @@ const buildBucket = ( }; const vectorize = {} as VectorizeIndex; +const d1 = {} as D1Database; beforeEach(() => { - vi.mocked(deleteKnowledgeBySource) + vi.mocked(indexKnowledgeSource) .mockReset() - .mockResolvedValue({ deleted: 0 }); - vi.mocked(processKnowledgeFile).mockReset().mockResolvedValue({ chunks: 2 }); + .mockResolvedValue({ indexed: true, status: "approved", chunks: 2 }); }); describe("syncAll", () => { it("md ファイルが無ければ空結果", async () => { const bucket = buildBucket([], {}); - const result = await syncAll({ bucket, vectorize, apiKey: "k" }); + const result = await syncAll({ bucket, vectorize, apiKey: "k", d1 }); expect(result).toEqual({ results: [], totalFiles: 0, totalChunks: 0, editedCount: 0, + skippedCount: 0, }); }); @@ -66,24 +65,22 @@ describe("syncAll", () => { buildR2Object("originals/photo.png"), ]; const bucket = buildBucket(objects, {}); - const result = await syncAll({ bucket, vectorize, apiKey: "k" }); + const result = await syncAll({ bucket, vectorize, apiKey: "k", d1 }); expect(result.totalFiles).toBe(0); - expect(processKnowledgeFile).not.toHaveBeenCalled(); + expect(indexKnowledgeSource).not.toHaveBeenCalled(); }); - it("md ファイルを 1 件処理: 削除 → embedding", async () => { + it("md ファイルを indexKnowledgeSource へ渡す", async () => { const objects = [buildR2Object("doc.md")]; const bucket = buildBucket(objects, { "doc.md": "# hello" }); - const result = await syncAll({ bucket, vectorize, apiKey: "k" }); + const result = await syncAll({ bucket, vectorize, apiKey: "k", d1 }); - expect(deleteKnowledgeBySource).toHaveBeenCalledWith(vectorize, "doc.md"); - expect(processKnowledgeFile).toHaveBeenCalledWith( + expect(indexKnowledgeSource).toHaveBeenCalledWith( "doc.md", "# hello", - vectorize, - "k", - undefined, + { d1, vectorize, apiKey: "k" }, + { r2Etag: "etag-doc.md" }, ); expect(result.totalFiles).toBe(1); expect(result.totalChunks).toBe(2); @@ -94,6 +91,21 @@ describe("syncAll", () => { }); }); + it("未承認で skip された情報源を skippedCount に数える", async () => { + vi.mocked(indexKnowledgeSource).mockResolvedValueOnce({ + indexed: false, + status: "pending", + chunks: 0, + }); + const objects = [buildR2Object("doc.md")]; + const bucket = buildBucket(objects, { "doc.md": "x" }); + + const result = await syncAll({ bucket, vectorize, apiKey: "k", d1 }); + + expect(result.skippedCount).toBe(1); + expect(result.results[0]).toMatchObject({ file: "doc.md", skipped: true }); + }); + it("bucket.get が null を返したら error 行を記録して続行", async () => { const objects = [buildR2Object("ghost.md"), buildR2Object("real.md")]; const bucket = buildBucket(objects, { @@ -101,7 +113,7 @@ describe("syncAll", () => { "real.md": "# r", }); - const result = await syncAll({ bucket, vectorize, apiKey: "k" }); + const result = await syncAll({ bucket, vectorize, apiKey: "k", d1 }); expect(result.totalFiles).toBe(2); expect(result.results.find((r) => r.file === "ghost.md")).toMatchObject({ @@ -118,7 +130,7 @@ describe("syncAll", () => { ]; const bucket = buildBucket(objects, { "doc.md": "edited content" }); - const result = await syncAll({ bucket, vectorize, apiKey: "k" }); + const result = await syncAll({ bucket, vectorize, apiKey: "k", d1 }); expect(result.editedCount).toBe(1); expect(result.results[0]?.edited).toBe(true); @@ -131,21 +143,23 @@ describe("syncAll", () => { ]; const bucket = buildBucket(objects, { "doc.md": "x" }); - const result = await syncAll({ bucket, vectorize, apiKey: "k" }); + const result = await syncAll({ bucket, vectorize, apiKey: "k", d1 }); expect(result.editedCount).toBe(0); expect(result.results[0]?.edited).toBe(false); }); - it("processKnowledgeFile のエラーを保持する", async () => { - vi.mocked(processKnowledgeFile).mockResolvedValueOnce({ + it("index のエラーを保持する", async () => { + vi.mocked(indexKnowledgeSource).mockResolvedValueOnce({ + indexed: true, + status: "approved", chunks: 0, error: "embed-failed", }); const objects = [buildR2Object("doc.md")]; const bucket = buildBucket(objects, { "doc.md": "x" }); - const result = await syncAll({ bucket, vectorize, apiKey: "k" }); + const result = await syncAll({ bucket, vectorize, apiKey: "k", d1 }); expect(result.results[0]).toMatchObject({ file: "doc.md", @@ -156,42 +170,53 @@ describe("syncAll", () => { }); it("複数 md を合計してチャンク数を返す", async () => { - vi.mocked(processKnowledgeFile) - .mockResolvedValueOnce({ chunks: 3 }) - .mockResolvedValueOnce({ chunks: 5 }); + vi.mocked(indexKnowledgeSource) + .mockResolvedValueOnce({ indexed: true, status: "approved", chunks: 3 }) + .mockResolvedValueOnce({ indexed: true, status: "approved", chunks: 5 }); const objects = [buildR2Object("a.md"), buildR2Object("b.md")]; const bucket = buildBucket(objects, { "a.md": "a", "b.md": "b" }); - const result = await syncAll({ bucket, vectorize, apiKey: "k" }); + const result = await syncAll({ bucket, vectorize, apiKey: "k", d1 }); expect(result.totalFiles).toBe(2); expect(result.totalChunks).toBe(8); }); }); describe("syncFile", () => { - it("delete → process を順に呼び結果を返す", async () => { - vi.mocked(processKnowledgeFile).mockResolvedValueOnce({ chunks: 7 }); + it("approveAs 付きで indexKnowledgeSource に委譲する", async () => { + vi.mocked(indexKnowledgeSource).mockResolvedValueOnce({ + indexed: true, + status: "approved", + chunks: 7, + }); const result = await syncFile("doc.md", "# c", { vectorize, apiKey: "k", + d1, + approveAs: "admin-1", }); - expect(deleteKnowledgeBySource).toHaveBeenCalledWith(vectorize, "doc.md"); - expect(processKnowledgeFile).toHaveBeenCalledWith( + expect(indexKnowledgeSource).toHaveBeenCalledWith( "doc.md", "# c", - vectorize, - "k", - undefined, + { d1, vectorize, apiKey: "k" }, + { approveAs: "admin-1" }, ); - expect(result).toEqual({ chunks: 7 }); + expect(result).toMatchObject({ chunks: 7, indexed: true }); }); - it("embedding が失敗したら error を伝播", async () => { - vi.mocked(processKnowledgeFile).mockResolvedValueOnce({ + it("index の error を伝播する", async () => { + vi.mocked(indexKnowledgeSource).mockResolvedValueOnce({ + indexed: true, + status: "approved", chunks: 0, error: "no-api-key", }); - const result = await syncFile("x.md", "y", { vectorize, apiKey: "k" }); + const result = await syncFile("x.md", "y", { + vectorize, + apiKey: "k", + d1, + approveAs: "admin-1", + }); expect(result.error).toBe("no-api-key"); }); }); diff --git a/server/src/services/knowledge/sync.ts b/server/src/services/knowledge/sync.ts index 9fca16c7..c8a00260 100644 --- a/server/src/services/knowledge/sync.ts +++ b/server/src/services/knowledge/sync.ts @@ -1,5 +1,5 @@ import { logger } from "~/lib/logger"; -import { deleteKnowledgeBySource, processKnowledgeFile } from "./embedding"; +import { indexKnowledgeSource } from "./indexing"; import { buildOriginalsMap, EDIT_THRESHOLD_MS } from "./utils"; type SyncResult = { @@ -7,6 +7,7 @@ type SyncResult = { chunks: number; error?: string; edited?: boolean; + skipped?: boolean; }; type SyncAllResult = { @@ -14,13 +15,14 @@ type SyncAllResult = { totalFiles: number; totalChunks: number; editedCount: number; + skippedCount: number; }; type SyncDeps = { bucket: R2Bucket; vectorize: VectorizeIndex; apiKey: string; - d1?: D1Database; + d1: D1Database; }; const isFileEdited = (mdFile: R2Object, originalsMap: Map) => { @@ -32,6 +34,16 @@ const isFileEdited = (mdFile: R2Object, originalsMap: Map) => { ); }; +export const listMarkdownObjects = async (bucket: R2Bucket) => { + const listed = await bucket.list(); + return { + allObjects: listed.objects, + mdFiles: listed.objects.filter( + (obj) => obj.key.endsWith(".md") && !obj.key.startsWith("originals/"), + ), + }; +}; + /** * R2バケットの全Markdownファイルを読み込み、Vectorizeに同期 */ @@ -41,12 +53,7 @@ export const syncAll = async ({ apiKey, d1, }: SyncDeps): Promise => { - const listed = await bucket.list(); - const allObjects = listed.objects; - - const mdFiles = allObjects.filter( - (obj) => obj.key.endsWith(".md") && !obj.key.startsWith("originals/"), - ); + const { allObjects, mdFiles } = await listMarkdownObjects(bucket); const originalsMap = buildOriginalsMap(allObjects); logger.info(`[Sync] Found ${mdFiles.length} markdown files`); @@ -66,13 +73,11 @@ export const syncAll = async ({ `[Sync] Processing ${obj.key} (${content.length} bytes)${edited ? " [EDITED]" : ""}`, ); - await deleteKnowledgeBySource(vectorize, obj.key); - const result = await processKnowledgeFile( + const result = await indexKnowledgeSource( obj.key, content, - vectorize, - apiKey, - d1, + { d1, vectorize, apiKey }, + { r2Etag: obj.etag }, ); results.push({ @@ -80,6 +85,7 @@ export const syncAll = async ({ chunks: result.chunks, error: result.error, edited, + skipped: result.indexed ? undefined : true, }); } @@ -88,6 +94,7 @@ export const syncAll = async ({ totalFiles: mdFiles.length, totalChunks: results.reduce((sum, r) => sum + r.chunks, 0), editedCount: results.filter((r) => r.edited).length, + skippedCount: results.filter((r) => r.skipped).length, }; }; @@ -97,14 +104,12 @@ export const syncAll = async ({ export const syncFile = async ( key: string, content: string, - deps: Omit, -): Promise<{ chunks: number; error?: string }> => { - await deleteKnowledgeBySource(deps.vectorize, key); - return processKnowledgeFile( + deps: Omit & { approveAs?: string }, +) => { + return indexKnowledgeSource( key, content, - deps.vectorize, - deps.apiKey, - deps.d1, + { d1: deps.d1, vectorize: deps.vectorize, apiKey: deps.apiKey }, + { approveAs: deps.approveAs }, ); }; diff --git a/server/src/services/knowledge/upload.test.ts b/server/src/services/knowledge/upload.test.ts index a83c814a..1f8520b9 100644 --- a/server/src/services/knowledge/upload.test.ts +++ b/server/src/services/knowledge/upload.test.ts @@ -42,7 +42,9 @@ const buildFile = ( }; beforeEach(() => { - vi.mocked(syncFile).mockReset().mockResolvedValue({ chunks: 3 }); + vi.mocked(syncFile) + .mockReset() + .mockResolvedValue({ indexed: true, status: "approved", chunks: 3 }); vi.mocked(convertToMarkdown).mockReset().mockResolvedValue("# converted"); vi.mocked(isSupportedMimeType).mockReset().mockReturnValue(true); }); @@ -56,6 +58,8 @@ describe("uploadMarkdownFile", () => { bucket, vectorize: {} as VectorizeIndex, apiKey: "k", + d1: {} as D1Database, + approveAs: "admin-1", }); expect(bucket.put).toHaveBeenCalledWith("doc.md", "# hello", { @@ -64,6 +68,8 @@ describe("uploadMarkdownFile", () => { expect(syncFile).toHaveBeenCalledWith("doc.md", "# hello", { vectorize: {}, apiKey: "k", + d1: {}, + approveAs: "admin-1", }); expect(result).toEqual({ key: "doc.md", chunks: 3 }); }); @@ -76,6 +82,8 @@ describe("uploadMarkdownFile", () => { bucket, vectorize: {} as VectorizeIndex, apiKey: "k", + d1: {} as D1Database, + approveAs: "admin-1", }); expect(result.key).toBe("custom.md"); @@ -88,6 +96,8 @@ describe("uploadMarkdownFile", () => { bucket, vectorize: {} as VectorizeIndex, apiKey: "k", + d1: {} as D1Database, + approveAs: "admin-1", }); expect(result.key).toBe("ready.md"); }); @@ -104,19 +114,28 @@ describe("uploadMarkdownFile", () => { bucket, vectorize: {} as VectorizeIndex, apiKey: "k", + d1: {} as D1Database, + approveAs: "admin-1", }), ).rejects.toThrow(/exceeds limit/); expect(bucket.put).not.toHaveBeenCalled(); }); it("syncFile の error を伝播", async () => { - vi.mocked(syncFile).mockResolvedValueOnce({ chunks: 0, error: "boom" }); + vi.mocked(syncFile).mockResolvedValueOnce({ + indexed: true, + status: "approved", + chunks: 0, + error: "boom", + }); const bucket = buildBucket(); const file = buildFile("x", { name: "doc.md" }); const result = await uploadMarkdownFile(file, null, { bucket, vectorize: {} as VectorizeIndex, apiKey: "k", + d1: {} as D1Database, + approveAs: "admin-1", }); expect(result.error).toBe("boom"); }); @@ -134,6 +153,8 @@ describe("convertAndUpload", () => { bucket, vectorize: {} as VectorizeIndex, apiKey: "k", + d1: {} as D1Database, + approveAs: "admin-1", }); // originals/photo.png に元ファイル @@ -165,6 +186,8 @@ describe("convertAndUpload", () => { bucket: buildBucket(), vectorize: {} as VectorizeIndex, apiKey: "k", + d1: {} as D1Database, + approveAs: "admin-1", }), ).rejects.toThrow(/exceeds limit/); }); @@ -177,6 +200,8 @@ describe("convertAndUpload", () => { bucket: buildBucket(), vectorize: {} as VectorizeIndex, apiKey: "k", + d1: {} as D1Database, + approveAs: "admin-1", }), ).rejects.toThrow(/Unsupported file type/); }); @@ -188,6 +213,8 @@ describe("convertAndUpload", () => { bucket, vectorize: {} as VectorizeIndex, apiKey: "k", + d1: {} as D1Database, + approveAs: "admin-1", }); // file.name に "." が無いので extension は "bin"... wait, "no-ext".split(".") = ["no-ext"], pop = "no-ext" // 実装的には pop || "bin" で "no-ext" になる @@ -208,6 +235,8 @@ describe("reconvertFromOriginal", () => { bucket, vectorize: {} as VectorizeIndex, apiKey: "k", + d1: {} as D1Database, + approveAs: "admin-1", }), ).rejects.toThrow(/Original file not found/); }); @@ -225,6 +254,8 @@ describe("reconvertFromOriginal", () => { bucket, vectorize: {} as VectorizeIndex, apiKey: "k", + d1: {} as D1Database, + approveAs: "admin-1", }), ).rejects.toThrow(/Unsupported file type/); }); @@ -240,6 +271,8 @@ describe("reconvertFromOriginal", () => { bucket, vectorize: {} as VectorizeIndex, apiKey: "k", + d1: {} as D1Database, + approveAs: "admin-1", }); expect(bucket.put).toHaveBeenCalledWith("x.md", "# converted", { @@ -265,6 +298,8 @@ describe("reconvertFromOriginal", () => { bucket, vectorize: {} as VectorizeIndex, apiKey: "k", + d1: {} as D1Database, + approveAs: "admin-1", }), ).rejects.toThrow(); expect(isSupportedMimeType).toHaveBeenCalledWith( diff --git a/server/src/services/knowledge/upload.ts b/server/src/services/knowledge/upload.ts index 55f6a5d9..d21ef200 100644 --- a/server/src/services/knowledge/upload.ts +++ b/server/src/services/knowledge/upload.ts @@ -9,7 +9,8 @@ type UploadDeps = { bucket: R2Bucket; vectorize: VectorizeIndex; apiKey: string; - d1?: D1Database; + d1: D1Database; + approveAs: string; }; type UploadResult = { @@ -55,6 +56,7 @@ export const uploadMarkdownFile = async ( vectorize: deps.vectorize, apiKey: deps.apiKey, d1: deps.d1, + approveAs: deps.approveAs, }); return { key, chunks: result.chunks, error: result.error }; @@ -110,6 +112,7 @@ export const convertAndUpload = async ( vectorize: deps.vectorize, apiKey: deps.apiKey, d1: deps.d1, + approveAs: deps.approveAs, }); return { @@ -159,6 +162,7 @@ export const reconvertFromOriginal = async ( vectorize: deps.vectorize, apiKey: deps.apiKey, d1: deps.d1, + approveAs: deps.approveAs, }); return { From a755ae713e35a9fb4bae42def94359e34b23b0bd Mon Sep 17 00:00:00 2001 From: owk-owk130 Date: Tue, 1 Sep 2026 17:11:26 +0900 Subject: [PATCH 6/6] =?UTF-8?q?feat(server):=20=E6=83=85=E5=A0=B1=E6=BA=90?= =?UTF-8?q?=E3=81=AE=E6=89=BF=E8=AA=8D=20API=20=E3=82=92=E8=BF=BD=E5=8A=A0?= =?UTF-8?q?=E3=81=99=E3=82=8B?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-Authored-By: Claude Fable 5 Claude-Session: https://claude.ai/code/session_018BjNFvQjDzMwF8HaZHnJjP --- server/src/routes/admin/knowledge/index.ts | 2 + .../routes/admin/knowledge/sources.test.ts | 329 ++++++++++++++++++ server/src/routes/admin/knowledge/sources.ts | 253 ++++++++++++++ 3 files changed, 584 insertions(+) create mode 100644 server/src/routes/admin/knowledge/sources.test.ts create mode 100644 server/src/routes/admin/knowledge/sources.ts diff --git a/server/src/routes/admin/knowledge/index.ts b/server/src/routes/admin/knowledge/index.ts index b18b8666..8e0e2ee5 100644 --- a/server/src/routes/admin/knowledge/index.ts +++ b/server/src/routes/admin/knowledge/index.ts @@ -3,6 +3,7 @@ import type { PrincipalVariables } from "~/lib/principal"; import { requireRole } from "~/middleware/require-role"; import { knowledgeConvertRoutes } from "./convert"; import { knowledgeFilesRoutes } from "./files"; +import { knowledgeSourcesRoutes } from "./sources"; import { knowledgeSyncRoutes } from "./sync"; export const knowledgeAdminRoutes = new OpenAPIHono<{ @@ -13,5 +14,6 @@ export const knowledgeAdminRoutes = new OpenAPIHono<{ knowledgeAdminRoutes.use("*", requireRole("super_admin")); knowledgeAdminRoutes.route("/", knowledgeSyncRoutes); +knowledgeAdminRoutes.route("/", knowledgeSourcesRoutes); knowledgeAdminRoutes.route("/", knowledgeFilesRoutes); knowledgeAdminRoutes.route("/", knowledgeConvertRoutes); diff --git a/server/src/routes/admin/knowledge/sources.test.ts b/server/src/routes/admin/knowledge/sources.test.ts new file mode 100644 index 00000000..8a20dc4f --- /dev/null +++ b/server/src/routes/admin/knowledge/sources.test.ts @@ -0,0 +1,329 @@ +import { beforeEach, describe, expect, it, vi } from "vitest"; + +vi.mock("~/repository/knowledge-source-repository", () => ({ + APPROVAL_STATUSES: ["pending", "approved", "rejected", "disabled"], + knowledgeSourceRepository: { + findByPath: vi.fn(), + list: vi.fn(), + insert: vi.fn(), + update: vi.fn(), + }, +})); + +vi.mock("~/services/knowledge/indexing", () => ({ + indexKnowledgeSource: vi.fn(), + removeKnowledgeSource: vi.fn(), +})); + +vi.mock("~/services/knowledge/sync", () => ({ + listMarkdownObjects: vi.fn(), +})); + +vi.mock("~/repository/admin-session-repository", () => ({ + adminSessionRepository: { findValid: vi.fn() }, +})); + +vi.mock("~/repository/admin-user-repository", () => ({ + adminUserRepository: { findById: vi.fn() }, +})); + +vi.mock("~/services/auth/anonymous-session", () => ({ + verifyAnonymousToken: vi.fn(), +})); + +const { knowledgeSourceRepository } = await import( + "~/repository/knowledge-source-repository" +); +const { indexKnowledgeSource, removeKnowledgeSource } = await import( + "~/services/knowledge/indexing" +); +const { listMarkdownObjects } = await import("~/services/knowledge/sync"); +const { adminSessionRepository } = await import( + "~/repository/admin-session-repository" +); +const { adminUserRepository } = await import( + "~/repository/admin-user-repository" +); +const { knowledgeAdminRoutes } = await import("."); + +import { withResolvePrincipal } from "~/__tests__/helpers/test-app"; + +const app = await withResolvePrincipal(knowledgeAdminRoutes); + +const testUser = { + id: "user-1", + username: "admin01", + name: "管理者", + role: "super_admin", + passwordHash: "100000:salt:hash", + createdAt: "2024-01-01T00:00:00Z", + updatedAt: null, +}; + +const r2Bucket = { get: vi.fn(), list: vi.fn() }; + +const mockEnv = { + DB: {} as D1Database, + KNOWLEDGE_BUCKET: r2Bucket as unknown as R2Bucket, + VECTORIZE: {} as VectorizeIndex, + GOOGLE_GENERATIVE_AI_API_KEY: "test-api-key", + JWT_SECRET: "test-secret-32-chars-long-enough", +} as unknown as CloudflareBindings; + +const VALID_OPAQUE_TOKEN = "a".repeat(64); + +const authedRequest = (path: string, init?: RequestInit) => { + const req = new Request(`http://localhost${path}`, init); + req.headers.set("Authorization", `Bearer ${VALID_OPAQUE_TOKEN}`); + return req; +}; + +const baseRow = { + sourcePath: "bus/index.md", + canonicalUrl: "https://example.com/bus", + sourceType: null, + sourceAuthority: null, + sourceHash: "hash-1", + r2Etag: "etag-1", + chunkCount: 5, + approvalStatus: "pending", + approvedBy: null, + approvedAt: null, + disabledAt: null, + verifiedAt: null, + indexedAt: null, + createdAt: "2026-09-01T00:00:00.000Z", + updatedAt: null, +}; + +beforeEach(() => { + vi.clearAllMocks(); + vi.mocked(adminSessionRepository.findValid).mockResolvedValue({ + token: VALID_OPAQUE_TOKEN, + userId: "user-1", + expiresAt: new Date(Date.now() + 86400000).toISOString(), + createdAt: "2024-01-01T00:00:00Z", + }); + vi.mocked(adminUserRepository.findById).mockResolvedValue(testUser); +}); + +describe("GET /sources", () => { + it("認証なしは 401", async () => { + const res = await app.request( + new Request("http://localhost/sources"), + undefined, + mockEnv, + ); + expect(res.status).toBe(401); + }); + + it("情報源一覧を返す", async () => { + vi.mocked(knowledgeSourceRepository.list).mockResolvedValue([baseRow]); + + const res = await app.request( + authedRequest("/sources"), + undefined, + mockEnv, + ); + + expect(res.status).toBe(200); + const body = (await res.json()) as { sources: unknown[] }; + expect(body.sources).toHaveLength(1); + expect(body.sources[0]).toMatchObject({ + sourcePath: "bus/index.md", + approvalStatus: "pending", + chunkCount: 5, + }); + }); +}); + +describe("POST /sources/backfill", () => { + it("未登録の情報源だけを approved で登録する", async () => { + vi.mocked(listMarkdownObjects).mockResolvedValue({ + allObjects: [], + mdFiles: [ + { key: "a.md", etag: "e-a" }, + { key: "b.md", etag: "e-b" }, + ] as R2Object[], + }); + vi.mocked(knowledgeSourceRepository.list).mockResolvedValue([ + { ...baseRow, sourcePath: "a.md" }, + ]); + r2Bucket.get.mockResolvedValue({ + text: async () => "---\nurl: 'https://example.com/b'\n---\n本文", + }); + + const res = await app.request( + authedRequest("/sources/backfill", { method: "POST" }), + undefined, + mockEnv, + ); + + expect(res.status).toBe(200); + const body = (await res.json()) as { registered: number; skipped: number }; + expect(body).toMatchObject({ registered: 1, skipped: 1 }); + expect(knowledgeSourceRepository.insert).toHaveBeenCalledTimes(1); + expect(knowledgeSourceRepository.insert).toHaveBeenCalledWith( + mockEnv.DB, + expect.objectContaining({ + sourcePath: "b.md", + canonicalUrl: "https://example.com/b", + approvalStatus: "approved", + approvedBy: "user-1", + r2Etag: "e-b", + }), + ); + }); +}); + +describe("PATCH /sources/status", () => { + const patchBody = (data: Record): RequestInit => ({ + method: "PATCH", + headers: { "Content-Type": "application/json" }, + body: JSON.stringify(data), + }); + + it("未登録の source_path は 404", async () => { + vi.mocked(knowledgeSourceRepository.findByPath).mockResolvedValue(null); + + const res = await app.request( + authedRequest( + "/sources/status", + patchBody({ sourcePath: "missing.md", action: "approve" }), + ), + undefined, + mockEnv, + ); + expect(res.status).toBe(404); + }); + + it("approve は R2 にファイルが無ければ 404", async () => { + vi.mocked(knowledgeSourceRepository.findByPath).mockResolvedValue(baseRow); + r2Bucket.get.mockResolvedValue(null); + + const res = await app.request( + authedRequest( + "/sources/status", + patchBody({ sourcePath: "bus/index.md", action: "approve" }), + ), + undefined, + mockEnv, + ); + expect(res.status).toBe(404); + }); + + it("approve は承認情報を記録して再インデックスする", async () => { + vi.mocked(knowledgeSourceRepository.findByPath) + .mockResolvedValueOnce(baseRow) + .mockResolvedValueOnce({ + ...baseRow, + approvalStatus: "approved", + approvedBy: "user-1", + }); + r2Bucket.get.mockResolvedValue({ text: async () => "# 本文" }); + vi.mocked(indexKnowledgeSource).mockResolvedValue({ + indexed: true, + status: "approved", + chunks: 4, + }); + + const res = await app.request( + authedRequest( + "/sources/status", + patchBody({ sourcePath: "bus/index.md", action: "approve" }), + ), + undefined, + mockEnv, + ); + + expect(res.status).toBe(200); + expect(knowledgeSourceRepository.update).toHaveBeenCalledWith( + mockEnv.DB, + "bus/index.md", + expect.objectContaining({ + approvalStatus: "approved", + approvedBy: "user-1", + disabledAt: null, + }), + ); + expect(indexKnowledgeSource).toHaveBeenCalledWith( + "bus/index.md", + "# 本文", + { + d1: mockEnv.DB, + vectorize: mockEnv.VECTORIZE, + apiKey: "test-api-key", + }, + ); + }); + + it("approve でインデックスが失敗したら 500", async () => { + vi.mocked(knowledgeSourceRepository.findByPath).mockResolvedValue(baseRow); + r2Bucket.get.mockResolvedValue({ text: async () => "# 本文" }); + vi.mocked(indexKnowledgeSource).mockResolvedValue({ + indexed: true, + status: "approved", + chunks: 0, + error: "embed failed", + }); + + const res = await app.request( + authedRequest( + "/sources/status", + patchBody({ sourcePath: "bus/index.md", action: "approve" }), + ), + undefined, + mockEnv, + ); + expect(res.status).toBe(500); + }); + + it.each([ + ["reject", "rejected"], + ["disable", "disabled"], + ] as const)("%s は検索対象からも削除する", async (action, status) => { + vi.mocked(knowledgeSourceRepository.findByPath) + .mockResolvedValueOnce({ ...baseRow, approvalStatus: "approved" }) + .mockResolvedValueOnce({ ...baseRow, approvalStatus: status }); + vi.mocked(removeKnowledgeSource).mockResolvedValue({ deleted: 5 }); + + const res = await app.request( + authedRequest( + "/sources/status", + patchBody({ sourcePath: "bus/index.md", action }), + ), + undefined, + mockEnv, + ); + + expect(res.status).toBe(200); + expect(knowledgeSourceRepository.update).toHaveBeenCalledWith( + mockEnv.DB, + "bus/index.md", + expect.objectContaining({ approvalStatus: status }), + ); + expect(removeKnowledgeSource).toHaveBeenCalledWith("bus/index.md", { + d1: mockEnv.DB, + vectorize: mockEnv.VECTORIZE, + }); + }); + + it("disable は disabledAt を記録する", async () => { + vi.mocked(knowledgeSourceRepository.findByPath) + .mockResolvedValueOnce({ ...baseRow, approvalStatus: "approved" }) + .mockResolvedValueOnce({ ...baseRow, approvalStatus: "disabled" }); + vi.mocked(removeKnowledgeSource).mockResolvedValue({ deleted: 5 }); + + await app.request( + authedRequest( + "/sources/status", + patchBody({ sourcePath: "bus/index.md", action: "disable" }), + ), + undefined, + mockEnv, + ); + + const input = vi.mocked(knowledgeSourceRepository.update).mock.calls[0][2]; + expect(input.disabledAt).toEqual(expect.any(String)); + }); +}); diff --git a/server/src/routes/admin/knowledge/sources.ts b/server/src/routes/admin/knowledge/sources.ts new file mode 100644 index 00000000..ec8ad800 --- /dev/null +++ b/server/src/routes/admin/knowledge/sources.ts @@ -0,0 +1,253 @@ +import { createRoute, OpenAPIHono, z } from "@hono/zod-openapi"; +import { HTTPException } from "hono/http-exception"; + +import { errorResponse } from "~/lib/openapi-errors"; +import type { PrincipalVariables } from "~/lib/principal"; +import { requireAdminUser } from "~/lib/principal"; +import { + APPROVAL_STATUSES, + type ApprovalStatus, + type KnowledgeSource, + knowledgeSourceRepository, +} from "~/repository/knowledge-source-repository"; +import { + indexKnowledgeSource, + removeKnowledgeSource, +} from "~/services/knowledge/indexing"; +import { buildSourceRecord } from "~/services/knowledge/source-meta"; +import { listMarkdownObjects } from "~/services/knowledge/sync"; +import { requireApiKey, validateFileKey } from "./schemas"; + +export const knowledgeSourcesRoutes = new OpenAPIHono<{ + Bindings: CloudflareBindings; + Variables: Partial; +}>(); + +const SourceSchema = z.object({ + sourcePath: z.string(), + canonicalUrl: z.string().nullable(), + sourceType: z.string().nullable(), + sourceAuthority: z.number().nullable(), + approvalStatus: z.enum(APPROVAL_STATUSES), + chunkCount: z.number(), + approvedBy: z.string().nullable(), + approvedAt: z.string().nullable(), + disabledAt: z.string().nullable(), + verifiedAt: z.string().nullable(), + indexedAt: z.string().nullable(), + createdAt: z.string(), + updatedAt: z.string().nullable(), +}); + +const toSourceResponse = (row: KnowledgeSource) => ({ + sourcePath: row.sourcePath, + canonicalUrl: row.canonicalUrl, + sourceType: row.sourceType, + sourceAuthority: row.sourceAuthority, + approvalStatus: row.approvalStatus as ApprovalStatus, + chunkCount: row.chunkCount, + approvedBy: row.approvedBy, + approvedAt: row.approvedAt, + disabledAt: row.disabledAt, + verifiedAt: row.verifiedAt, + indexedAt: row.indexedAt, + createdAt: row.createdAt, + updatedAt: row.updatedAt, +}); + +const listSourcesRoute = createRoute({ + method: "get", + path: "/sources", + summary: "情報源一覧を取得", + description: "ナレッジ情報源の承認状態を一覧します", + tags: ["Admin - Knowledge"], + responses: { + 200: { + description: "情報源一覧", + content: { + "application/json": { + schema: z.object({ sources: z.array(SourceSchema) }), + }, + }, + }, + 401: errorResponse(401), + 403: errorResponse(403), + }, +}); + +knowledgeSourcesRoutes.openapi(listSourcesRoute, async (c) => { + const rows = await knowledgeSourceRepository.list(c.env.DB); + return c.json({ sources: rows.map(toSourceResponse) }, 200); +}); + +const backfillRoute = createRoute({ + method: "post", + path: "/sources/backfill", + summary: "既存ナレッジを承認済み情報源として登録", + description: + "R2 の Markdown のうち情報源未登録のものを approved で一括登録します", + tags: ["Admin - Knowledge"], + responses: { + 200: { + description: "登録結果", + content: { + "application/json": { + schema: z.object({ + message: z.string(), + registered: z.number(), + skipped: z.number(), + }), + }, + }, + }, + 401: errorResponse(401), + 403: errorResponse(403), + }, +}); + +const BACKFILL_BATCH_SIZE = 20; + +knowledgeSourcesRoutes.openapi(backfillRoute, async (c) => { + const adminUser = requireAdminUser(c.get("principal")); + const { mdFiles } = await listMarkdownObjects(c.env.KNOWLEDGE_BUCKET); + const registeredPaths = new Set( + (await knowledgeSourceRepository.list(c.env.DB)).map((r) => r.sourcePath), + ); + const targets = mdFiles.filter((obj) => !registeredPaths.has(obj.key)); + const now = new Date().toISOString(); + + let registered = 0; + for (let i = 0; i < targets.length; i += BACKFILL_BATCH_SIZE) { + const batch = targets.slice(i, i + BACKFILL_BATCH_SIZE); + const results = await Promise.all( + batch.map(async (obj) => { + const file = await c.env.KNOWLEDGE_BUCKET.get(obj.key); + if (!file) return false; + await knowledgeSourceRepository.insert(c.env.DB, { + ...(await buildSourceRecord(obj.key, await file.text())), + r2Etag: obj.etag, + approvalStatus: "approved", + approvedBy: adminUser.id, + approvedAt: now, + createdAt: now, + }); + return true; + }), + ); + registered += results.filter(Boolean).length; + } + + return c.json( + { + message: `${registered}件の情報源を登録しました`, + registered, + skipped: mdFiles.length - registered, + }, + 200, + ); +}); + +const updateStatusRoute = createRoute({ + method: "patch", + path: "/sources/status", + summary: "情報源の承認状態を変更", + description: + "approve は再インデックス、reject / disable は検索対象からの削除まで行います", + tags: ["Admin - Knowledge"], + request: { + body: { + content: { + "application/json": { + schema: z.object({ + sourcePath: z.string().min(1), + action: z.enum(["approve", "reject", "disable"]), + }), + }, + }, + required: true, + }, + }, + responses: { + 200: { + description: "変更結果", + content: { + "application/json": { + schema: z.object({ + message: z.string(), + source: SourceSchema, + }), + }, + }, + }, + 401: errorResponse(401), + 403: errorResponse(403), + 404: errorResponse(404), + 500: errorResponse(500), + }, +}); + +knowledgeSourcesRoutes.openapi(updateStatusRoute, async (c) => { + const { sourcePath, action } = c.req.valid("json"); + validateFileKey(sourcePath); + const adminUser = requireAdminUser(c.get("principal")); + + const row = await knowledgeSourceRepository.findByPath(c.env.DB, sourcePath); + if (!row) { + throw new HTTPException(404, { message: "情報源が見つかりません" }); + } + + const now = new Date().toISOString(); + + if (action === "approve") { + const file = await c.env.KNOWLEDGE_BUCKET.get(sourcePath); + if (!file) { + throw new HTTPException(404, { + message: "R2 にファイルが見つかりません", + }); + } + const apiKey = requireApiKey(c.env.GOOGLE_GENERATIVE_AI_API_KEY); + + await knowledgeSourceRepository.update(c.env.DB, sourcePath, { + approvalStatus: "approved", + approvedBy: adminUser.id, + approvedAt: now, + disabledAt: null, + }); + + const result = await indexKnowledgeSource(sourcePath, await file.text(), { + d1: c.env.DB, + vectorize: c.env.VECTORIZE, + apiKey, + }); + if (result.error) { + throw new HTTPException(500, { + message: `インデックスに失敗しました: ${result.error}`, + }); + } + } else { + await knowledgeSourceRepository.update(c.env.DB, sourcePath, { + approvalStatus: action === "reject" ? "rejected" : "disabled", + ...(action === "disable" && { disabledAt: now }), + }); + await removeKnowledgeSource(sourcePath, { + d1: c.env.DB, + vectorize: c.env.VECTORIZE, + }); + } + + const updated = await knowledgeSourceRepository.findByPath( + c.env.DB, + sourcePath, + ); + if (!updated) { + throw new HTTPException(500, { message: "情報源の更新に失敗しました" }); + } + + return c.json( + { + message: `${sourcePath} を ${updated.approvalStatus} にしました`, + source: toSourceResponse(updated), + }, + 200, + ); +});