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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
215 changes: 205 additions & 10 deletions scripts/packed-process-smoke.mjs
Original file line number Diff line number Diff line change
Expand Up @@ -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 { <https://packed.invalid/entity/Q1> <https://packed.invalid/prop/direct/P1> ?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 { <https://packed.invalid/entity/Q1> } 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;
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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,
Expand All @@ -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',
});
Expand Down Expand Up @@ -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',
Expand Down
33 changes: 16 additions & 17 deletions src/local-dependencies.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,4 @@
import { Buffer } from 'node:buffer';
import { Readable, Transform, type TransformCallback } from 'node:stream';
import {
createAuthorizedTaproot,
commitAuthorizationBootstrapV1,
Expand Down Expand Up @@ -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<RDF.Quad>;
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<RDF.Quad>;
) 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<RDF.Quad> & {
filter(predicate: (quad: RDF.Quad) => boolean): RDF.Stream<RDF.Quad>;
};

function scopedPolicy(value: unknown): {
kind: 'taproot-scoped-sparql-policy-v1';
datasetId: string;
Expand Down
10 changes: 9 additions & 1 deletion src/server/tasks.ts
Original file line number Diff line number Diff line change
Expand Up @@ -695,6 +695,14 @@ export class TaskService {
): Promise<Task> {
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();
Expand Down Expand Up @@ -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,
Expand Down
46 changes: 44 additions & 2 deletions test/search.integration.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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',
Expand Down Expand Up @@ -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,
Expand Down
Loading