diff --git a/knowledge-fs/packages/api/src/research-retrieval-checkpoint.test.ts b/knowledge-fs/packages/api/src/research-retrieval-checkpoint.test.ts index cbd5e5197f6..98817b6c9c7 100644 --- a/knowledge-fs/packages/api/src/research-retrieval-checkpoint.test.ts +++ b/knowledge-fs/packages/api/src/research-retrieval-checkpoint.test.ts @@ -235,6 +235,72 @@ describe("Research retrieval durable search checkpoint", () => { ).toThrow("checkpoint trace mismatch"); }); + it("scopes the durable checkpoint by the vision-expanded retrieval query", () => { + const retrievalQuery = "invoice retention\n\nImage OCR: invoice 42"; + const queryImage = { + byteSize: 3, + mimeType: "image/png" as const, + sha256: "a".repeat(64), + uploadFileId: "70000000-0000-4000-8000-000000000001", + }; + // The first V3 boundary Research persists: planned, no judgement, no tree frontier. + const searchState = (query: string) => ({ + budget: { + elapsedMs: 10, + exhaustedReasons: [], + modelCalls: 1, + openedResources: 0, + retrievalSteps: 0, + rounds: 0, + supplementalSearches: 0, + }, + fingerprint: `projection-set-sha256:${"b".repeat(64)}`, + knowledgeSpaceId: SPACE_ID, + phase: "planned" as const, + publicationId: PUBLICATION_ID, + query, + queryPlan: { + evidenceDimensions: ["retention"], + intent: "direct" as const, + subqueries: [], + useGraph: false, + }, + sequence: 0, + tenantId: "tenant-1", + traceId: TRACE_ID, + version: ResearchEvidenceRetrievalCheckpointVersion, + }); + + // A mixed text+image run keeps the user's text as `query` and retrieves with the expansion. + expect( + validateResearchRetrievalDurableCheckpoint({ + evidenceBundle: { ...evidenceBundle(), queryImages: [queryImage], retrievalQuery }, + searchState: searchState(retrievalQuery), + }).evidenceBundle, + ).toMatchObject({ query: "invoice retention", retrievalQuery }); + + // An image-only run has no user text at all. + expect( + validateResearchRetrievalDurableCheckpoint({ + evidenceBundle: { + ...evidenceBundle(), + query: "", + queryImages: [queryImage], + retrievalQuery: "Image OCR: invoice 42", + }, + searchState: searchState("Image OCR: invoice 42"), + }).searchState.query, + ).toBe("Image OCR: invoice 42"); + + // The search state must still match what was actually retrieved, not the raw text. + expect(() => + validateResearchRetrievalDurableCheckpoint({ + evidenceBundle: { ...evidenceBundle(), queryImages: [queryImage], retrievalQuery }, + searchState: searchState("invoice retention"), + }), + ).toThrow("durable checkpoint scope mismatch"); + }); + it("rehydrates checkpoint evidence with bounded citation and source fallbacks", () => { const base = evidenceBundle(); const item = { diff --git a/knowledge-fs/packages/api/src/research-retrieval-checkpoint.ts b/knowledge-fs/packages/api/src/research-retrieval-checkpoint.ts index c7f144138d6..31a0b5ab4d1 100644 --- a/knowledge-fs/packages/api/src/research-retrieval-checkpoint.ts +++ b/knowledge-fs/packages/api/src/research-retrieval-checkpoint.ts @@ -304,8 +304,12 @@ export function validateResearchRetrievalDurableCheckpoint( const envelope = durableCheckpointEnvelopeSchema.parse(value); const evidenceBundle = EvidenceBundleSchema.parse(envelope.evidenceBundle); const searchState = parseAnyResearchRetrievalSearchCheckpoint(envelope.searchState); + // The search state is scoped by the query that was actually retrieved. A query-image run keeps + // the user's text in `query` and records the vision-expanded text it retrieved with in + // `retrievalQuery`, so the scope must be compared against the latter when present. + const retrievedQuery = evidenceBundle.retrievalQuery ?? evidenceBundle.query; if ( - evidenceBundle.query !== searchState.query || + retrievedQuery !== searchState.query || (evidenceBundle.traceId !== undefined && evidenceBundle.traceId !== searchState.traceId) ) { throw new Error("Research retrieval durable checkpoint scope mismatch"); diff --git a/knowledge-fs/packages/api/src/research-task-runtime.test.ts b/knowledge-fs/packages/api/src/research-task-runtime.test.ts index 53583a94910..423ec733023 100644 --- a/knowledge-fs/packages/api/src/research-task-runtime.test.ts +++ b/knowledge-fs/packages/api/src/research-task-runtime.test.ts @@ -957,6 +957,119 @@ describe("research task production runtime", () => { }); }); + it("persists and resumes a text-plus-image Research checkpoint scoped by the expanded query", async () => { + const frozenRuntime = publishedRuntimeSnapshot(SPACE_ID); + const imageId = "00000000-0000-4000-8000-000000000001"; + const expansion = "Image OCR: invoice 42"; + const repository = new MemoryDurableRepository({ + ...baseJob(), + metadata: { + [QUERY_IMAGE_REFERENCES_METADATA_KEY]: [{ uploadFileId: imageId }], + [RESEARCH_TASK_RUNTIME_SNAPSHOT_METADATA_KEY]: + toResearchTaskRuntimeSnapshotPayload(frozenRuntime), + }, + mode: "research", + query: "What does this diagram show?", + }); + const partials = createInMemoryResearchTaskPartialResultRepository({ + maxListLimit: 10, + maxResults: 10, + }); + const generationInputs: Array> = []; + let generationAttempt = 0; + let now = 1_000; + const durableCheckpoint = (retrievalQuery: string) => ({ + evidenceBundle: { + ...evidenceBundle(), + query: "What does this diagram show?", + queryImages: [ + { + byteSize: 3, + mimeType: "image/png" as const, + sha256: "a".repeat(64), + uploadFileId: imageId, + }, + ], + retrievalQuery, + traceId: JOB_ID, + }, + searchState: { + ...durableRetrievalCheckpoint(frozenRuntime).searchState, + query: retrievalQuery, + }, + }); + const runtime = createResearchTaskRuntime({ + ...runtimeOptions(repository), + allowLegacyProfileFallback: false, + generator: { + stream: async function* (input) { + generationInputs.push(input as unknown as Record); + generationAttempt += 1; + // Mirrors createQueryImageAwareQueryGenerator: expand once, then retrieve with the + // user's text joined to the expansion while the bundle keeps the raw text as `query`. + const persistedExpansion = input.queryImageExpansion ?? expansion; + if (!input.queryImageExpansion) await input.onQueryImageExpansion?.(expansion); + const retrievalQuery = [input.query.trim(), persistedExpansion].join("\n\n"); + if (generationAttempt === 1) { + await input.onResearchDurableCheckpoint?.(durableCheckpoint(retrievalQuery)); + throw new Error("answer provider timed out after retrieval"); + } + expect(input.researchDurableCheckpoint).toEqual(durableCheckpoint(retrievalQuery)); + yield traceStep("query.retrieve", { checkpointed: true, itemCount: 1 }); + yield traceStep("query.answer"); + yield { delta: "Diagram answer", type: "delta" as const }; + yield { + finishReason: "retrieval-evidence", + metadata: { evidenceBundle: input.researchDurableCheckpoint?.evidenceBundle }, + type: "done" as const, + }; + }, + }, + maxRetryDelayMs: 1, + now: () => now, + partials, + projectionSnapshotResolver: { resolve: async () => frozenRuntime.projectionSnapshot }, + queryImageResolver: { + resolve: async () => [ + { + body: new Uint8Array([1, 2, 3]), + byteSize: 3, + mimeType: "image/png" as const, + sha256: "a".repeat(64), + uploadFileId: imageId, + }, + ], + }, + retryDelayMs: 1, + }); + + await expect(runtime.tick()).resolves.toMatchObject({ retryScheduled: 1, succeeded: 0 }); + expect(repository.job.error).toBe("answer provider timed out after retrieval"); + expect(repository.job.metadata).toHaveProperty( + RESEARCH_RETRIEVAL_DURABLE_CHECKPOINT_METADATA_KEY, + ); + expect(repository.job.metadata[QUERY_IMAGE_EXPANSION_METADATA_KEY]).toBe(expansion); + + now = 1_002; + await expect(runtime.tick()).resolves.toMatchObject({ retryScheduled: 0, succeeded: 1 }); + expect(repository.job.stage).toBe("completed"); + expect(generationInputs).toHaveLength(2); + expect(generationInputs[1]).toMatchObject({ + queryImageExpansion: expansion, + researchDurableCheckpoint: durableCheckpoint(`What does this diagram show?\n\n${expansion}`), + }); + // The retrieval boundary is streamed as a partial result before the answer is persisted. + const persisted = await partials.list({ + limit: 10, + researchTaskJobId: JOB_ID, + tenantId: "tenant-1", + }); + expect(persisted.items.at(-1)).toMatchObject({ + answer: "Diagram answer", + evidenceBundle: { query: "What does this diagram show?", retrievalQuery: expect.any(String) }, + }); + }); + it("reserves each Research model call and reconciles it with Dify token usage", async () => { const frozenRuntime = publishedRuntimeSnapshot(SPACE_ID); const repository = new MemoryDurableRepository({