From 992e9647c7cee5929447ebeb771c0083c894bed3 Mon Sep 17 00:00:00 2001 From: Kendrik Felty Date: Sat, 25 Jul 2026 00:28:08 -0500 Subject: [PATCH] fix: close archive and SPARQL boundaries --- scripts/packed-process-smoke.mjs | 215 +++++++++++++++++++++++++++++-- src/local-dependencies.ts | 33 +++-- src/server/tasks.ts | 10 +- test/search.integration.test.ts | 46 ++++++- 4 files changed, 274 insertions(+), 30 deletions(-) diff --git a/scripts/packed-process-smoke.mjs b/scripts/packed-process-smoke.mjs index 4e58086..14f8c1a 100644 --- a/scripts/packed-process-smoke.mjs +++ b/scripts/packed-process-smoke.mjs @@ -1169,6 +1169,92 @@ export async function createWorkshopSemanticConfigurationDependencies( ); assert.equal(populated.format, 'gnolith-workshop-public-backup-fixture-v1'); assert.equal(populated.q1Visibility.target.id, 'Q1'); + const globalTriplePatternQuery = + 'SELECT ?subject ?predicate ?object WHERE { ?subject ?predicate ?object } LIMIT 20'; + const validatedGlobalTriplePattern = await successfulToolCall( + origin, + token, + 166, + 'validate_sparql', + { query: globalTriplePatternQuery }, + ); + assert.equal(validatedGlobalTriplePattern.valid, true); + assert.equal(validatedGlobalTriplePattern.queryType, 'SELECT'); + const dryRunGlobalTriplePattern = await successfulToolCall( + origin, + token, + 167, + 'query_sparql', + { query: globalTriplePatternQuery, dryRun: true }, + ); + assert.equal(dryRunGlobalTriplePattern.valid, true); + assert.equal(dryRunGlobalTriplePattern.dryRun, true); + const globalTriplePattern = await successfulToolCall( + origin, + token, + 168, + 'query_sparql', + { query: globalTriplePatternQuery }, + ); + assert.equal(globalTriplePattern.status, 200); + assert.ok( + globalTriplePattern.data.results.bindings.length >= 1, + JSON.stringify(globalTriplePattern), + ); + const triplePatternQuery = + 'SELECT ?value WHERE { ?value }'; + const validatedTriplePattern = await successfulToolCall( + origin, + token, + 170, + 'validate_sparql', + { query: triplePatternQuery }, + ); + assert.equal(validatedTriplePattern.valid, true); + assert.equal(validatedTriplePattern.queryType, 'SELECT'); + const dryRunTriplePattern = await successfulToolCall( + origin, + token, + 171, + 'query_sparql', + { query: triplePatternQuery, dryRun: true }, + ); + assert.equal(dryRunTriplePattern.valid, true); + assert.equal(dryRunTriplePattern.dryRun, true); + assert.ok(dryRunTriplePattern.allowedSubjectCount >= 2); + const emptyGlobalQuery = await successfulToolCall( + origin, + token, + 172, + 'query_sparql', + { + query: + 'SELECT ?subject WHERE { VALUES ?subject { } FILTER(false) }', + }, + ); + assert.equal(emptyGlobalQuery.status, 200); + assert.deepEqual(emptyGlobalQuery.data.results.bindings, []); + const boundTriplePattern = await successfulToolCall( + origin, + token, + 173, + 'query_sparql', + { query: triplePatternQuery }, + ); + assert.equal(boundTriplePattern.status, 200); + assert.ok( + boundTriplePattern.data.results.bindings.length >= 1, + JSON.stringify(boundTriplePattern), + ); + const unauthorizedTriplePattern = await successfulToolCall( + origin, + outsiderToken, + 174, + 'query_sparql', + { query: triplePatternQuery }, + ); + assert.equal(unauthorizedTriplePattern.status, 200); + assert.deepEqual(unauthorizedTriplePattern.data.results.bindings, []); const q1PolicyRevision = ( await successfulToolCall(origin, token, 180, 'gnolith_status', {}) ).authorizationRevision; @@ -1474,10 +1560,73 @@ export async function createWorkshopSemanticConfigurationDependencies( prompt: 'Archive after the lifecycle projection is committed.', }, ); - await successfulToolCall(origin, token, 2312, 'archive_task', { - taskId: transientTask.id, - expectedRevision: transientTask.revision, + const transientTaskUpdated = await successfulToolCall( + origin, + token, + 2390, + 'update_task', + { + taskId: transientTask.id, + expectedRevision: transientTask.revision, + patch: { + title: 'Archive the revised packed mixed lifecycle task', + description: 'Archive the revised packed mixed lifecycle task', + }, + }, + ); + const transientTaskArchived = await successfulToolCall( + origin, + token, + 2312, + 'archive_task', + { + taskId: transientTask.id, + expectedRevision: transientTaskUpdated.revision, + }, + ); + assert.equal(transientTaskArchived.revision, 3); + assert.ok(transientTaskArchived.archivedAt); + const transientTaskHistory = await successfulToolCall( + origin, + token, + 2391, + 'task_history', + { taskId: transientTask.id, limit: 10 }, + ); + assert.deepEqual( + transientTaskHistory.items.map(({ revision }) => revision), + [3, 2, 1], + ); + const beforeArchiveReplay = await successfulToolCall( + origin, + token, + 2392, + 'search_admin', + { action: 'health' }, + ); + const archiveReplay = await mcp(origin, token, 2393, 'tools/call', { + name: 'archive_task', + arguments: { + taskId: transientTask.id, + expectedRevision: transientTaskArchived.revision, + }, }); + assert.equal( + toolErrorCode(archiveReplay), + 'conflict', + JSON.stringify(archiveReplay), + ); + const afterArchiveReplay = await successfulToolCall( + origin, + token, + 2394, + 'search_admin', + { action: 'health' }, + ); + assert.equal( + afterArchiveReplay.materialization.sourceHighWatermark, + beforeArchiveReplay.materialization.sourceHighWatermark, + ); await successfulToolCall(origin, token, 2313, 'delete_annotation', { annotationId: transientAnnotation.id, expectedRevision: transientAnnotationUpdated.revision, @@ -1511,13 +1660,11 @@ export async function createWorkshopSemanticConfigurationDependencies( mixedLagged.materialization.activeAppliedWatermark, JSON.stringify(mixedLagged), ); - for (let page = 0; page < 10; page += 1) { - await successfulToolCall(origin, token, 2318 + page, 'search_admin', { - action: 'materialize', - limit: 100, - options: { maxRebuildRoots: 100 }, - }); - } + await successfulToolCall(origin, token, 2318, 'search_admin', { + action: 'materialize', + limit: 100, + options: { maxRebuildRoots: 100 }, + }); const mixedCaughtUp = await successfulToolCall( origin, token, @@ -1542,6 +1689,24 @@ export async function createWorkshopSemanticConfigurationDependencies( activeAppliedWatermark: mixedCaughtUp.materialization.sourceHighWatermark, }, ); + const archivedTaskSearch = await successfulToolCall( + origin, + token, + 2395, + 'search', + { + text: 'Archive the revised packed mixed lifecycle task', + kinds: ['task'], + limit: 20, + }, + ); + assert.equal( + archivedTaskSearch.items.filter( + ({ sourceId }) => sourceId === transientTask.id, + ).length, + 0, + JSON.stringify(archivedTaskSearch), + ); await successfulToolCall(origin, token, 2329, 'search_admin', { action: 'rebuild', }); @@ -1680,6 +1845,36 @@ export async function createWorkshopSemanticConfigurationDependencies( await waitForLive(origin, child); const qdrantRestarted = await waitForReady(origin, token, child); assert.equal(qdrantRestarted.semanticState.ready, true); + const boundTriplePatternAfterRestart = await successfulToolCall( + origin, + token, + 175, + 'query_sparql', + { query: triplePatternQuery }, + ); + assert.equal(boundTriplePatternAfterRestart.status, 200); + assert.ok( + boundTriplePatternAfterRestart.data.results.bindings.length >= 1, + JSON.stringify(boundTriplePatternAfterRestart), + ); + const archivedTaskSearchAfterRestart = await successfulToolCall( + origin, + token, + 2396, + 'search', + { + text: 'Archive the revised packed mixed lifecycle task', + kinds: ['task'], + limit: 20, + }, + ); + assert.equal( + archivedTaskSearchAfterRestart.items.filter( + ({ sourceId }) => sourceId === transientTask.id, + ).length, + 0, + JSON.stringify(archivedTaskSearchAfterRestart), + ); assert.ok( [...semanticMockPoints.values()].some( (point) => point.payload?.sourceId === 'Q1', diff --git a/src/local-dependencies.ts b/src/local-dependencies.ts index 4d378b8..35ca1a2 100644 --- a/src/local-dependencies.ts +++ b/src/local-dependencies.ts @@ -1,5 +1,4 @@ import { Buffer } from 'node:buffer'; -import { Readable, Transform, type TransformCallback } from 'node:stream'; import { createAuthorizedTaproot, commitAuthorizationBootstrapV1, @@ -648,31 +647,31 @@ function scopedSource( object?: RDF.Term | null, graph?: RDF.Term | null, ) { - if (subject?.termType === 'NamedNode' && !allowed.has(subject.value)) - return Readable.from([], { - objectMode: true, - }) as RDF.Stream; - const filter = new Transform({ - objectMode: true, - transform( - quad: RDF.Quad, - _encoding: BufferEncoding, - callback: TransformCallback, - ) { - callback(null, allowed.has(quad.subject.value) ? quad : undefined); - }, - }); const matches = source.match( subject, predicate, object, graph, - ) as unknown as Readable; - return matches.pipe(filter) as RDF.Stream; + ) as unknown as FilterableRdfStream; + if (typeof matches.filter !== 'function') + throw new WorkshopError( + 'dependency_unavailable', + 'Diamond RDF source is incompatible', + 503, + ); + return matches.filter( + subject?.termType === 'NamedNode' && !allowed.has(subject.value) + ? () => false + : (quad) => allowed.has(quad.subject.value), + ); }, }; } +type FilterableRdfStream = RDF.Stream & { + filter(predicate: (quad: RDF.Quad) => boolean): RDF.Stream; +}; + function scopedPolicy(value: unknown): { kind: 'taproot-scoped-sparql-policy-v1'; datasetId: string; diff --git a/src/server/tasks.ts b/src/server/tasks.ts index 2a4b3a8..79d7acd 100644 --- a/src/server/tasks.ts +++ b/src/server/tasks.ts @@ -695,6 +695,14 @@ export class TaskService { ): Promise { const context = await this.#context(authorization, 'task:write'); const current = await this.#getForMutation(id, context); + if (current.archivedAt) + throw new WorkshopError('conflict', 'Task is already archived', 409); + if (current.completedAt) + throw new WorkshopError( + 'conflict', + 'Completed tasks cannot be archived', + 409, + ); const predicate = visibilitySql('workshop_tasks', context); const expected = revisionPrecondition(options); const now = this.#clock().toISOString(); @@ -1583,7 +1591,7 @@ export class TaskService { context, eventId, sourceId: next.id, - operation: 'upsert', + operation: next.archivedAt ? 'delete' : 'upsert', changeClass, sourceRevision: next.revision, sourcePolicyRevision: next.policyRevision, diff --git a/test/search.integration.test.ts b/test/search.integration.test.ts index 397881d..a216572 100644 --- a/test/search.integration.test.ts +++ b/test/search.integration.test.ts @@ -1180,9 +1180,41 @@ describe('Workshop external search producers', () => { { description: 'Archive producer changes', prompt: 'Archive safely.' }, context, ); + const revisedArchivedTask = await core.tasks.update( + archivedTask.id, + { + expectedRevision: archivedTask.revision, + title: 'Archive revised producer changes', + }, + context, + ); await expect( - core.tasks.archive(archivedTask.id, { expectedRevision: 1 }, context), - ).resolves.toMatchObject({ archivedAt: expect.any(String) }); + core.tasks.archive( + archivedTask.id, + { expectedRevision: revisedArchivedTask.revision }, + context, + ), + ).resolves.toMatchObject({ + revision: 3, + archivedAt: expect.any(String), + }); + expect( + (await core.tasks.history(archivedTask.id, context)).map( + ({ revision }) => revision, + ), + ).toEqual([3, 2, 1]); + const beforeArchiveReplay = await search.materialization.health(context); + await expect( + core.tasks.archive(archivedTask.id, { expectedRevision: 3 }, context), + ).rejects.toMatchObject({ + code: 'conflict', + message: 'Task is already archived', + }); + await expect( + search.materialization.health(context), + ).resolves.toMatchObject({ + sourceHighWatermark: beforeArchiveReplay.sourceHighWatermark, + }); const abandonedTask = await core.tasks.create( { description: 'Reset abandoned producer claim', @@ -1231,6 +1263,16 @@ describe('Workshop external search producers', () => { maxJobs: 20, maxRebuildRoots: 20, }); + await expect( + search.search( + { + text: 'Archive revised producer changes', + kinds: ['task'], + limit: 10, + }, + context, + ), + ).resolves.toMatchObject({ results: [] }); const first = await search.search( { text: 'overlapping work', kinds: ['prompt'], limit: 10 }, context,