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
26 changes: 4 additions & 22 deletions apps/server/src/storage/contract.ts
Original file line number Diff line number Diff line change
@@ -1,3 +1,5 @@
import { researchWorkspace } from "./research-inline-contract-fixtures";
import { researchInlineContract } from "./research-inline-contract";
import { describe, expect, it } from "bun:test";

import {
Expand All @@ -11,33 +13,12 @@ import { StorageError } from "./errors";
import { BACKGROUND_JOB_PROGRESS_LIMIT } from "./model";

import type { StorageFactory as Factory } from "./contract-support";
import type { CreateResearchWorkspace, JsonValue, Lease } from "./model";
import type { JsonValue } from "./model";

function attempt<T>(action: () => Promise<T>): Promise<T> {
return Promise.resolve().then(action);
}

function researchWorkspace(
channelId: string,
createdBy: string,
lease: Lease,
overrides: Partial<CreateResearchWorkspace> = {},
): CreateResearchWorkspace {
return {
id: id("research-workspace"),
channelId,
title: "API compatibility research",
proposedQuestion: "Which API contracts changed?",
origin: "sidebar",
createdBy,
idempotencyKey: id("create-research"),
fingerprint: id("research-fingerprint"),
now: new Date("2026-01-07T03:04:05.000Z"),
lease,
...overrides,
};
}

/** The behavioral gate every built-in storage adapter must pass. */
export function storageContract(name: string, factory: Factory): void {
describe(`${name} storage`, () => {
Expand Down Expand Up @@ -1932,6 +1913,7 @@ export function storageContract(name: string, factory: Factory): void {
});

researchPublicationContract(factory);
researchInlineContract(factory);

it("rejects new archived research mutations but preserves exact replays", async () => {
let storage = await opened(factory);
Expand Down
74 changes: 74 additions & 0 deletions apps/server/src/storage/memory/adapter.test.ts
Original file line number Diff line number Diff line change
@@ -1,4 +1,78 @@
import { expect, it } from "bun:test";
import { storageContract } from "../contract";
import { backgroundJob, contractId as id, userAndChannel } from "../contract-support";
import { MemoryStorage } from "./adapter";

storageContract("memory", () => new MemoryStorage());

it("commits all research projections before a linked job can become terminal", async () => {
let storage = new MemoryStorage();
await storage.migrate();
try {
let { channelId, userId, lease } = await userAndChannel(storage);
let now = new Date("2026-01-07T03:04:05.000Z");
let references: { id: string; jobId: string }[] = [];
for (let index = 0; index < 2; index++) {
let workspaceId = id(`projection-workspace-${index}`);
let started = await storage.research.start({
id: workspaceId,
channelId,
title: "Inline research",
question: `Question ${index}`,
origin: "inline",
createdBy: userId,
turnId: id(`projection-turn-${index}`),
messageId: id(`projection-message-${index}`),
requestId: id(`projection-request-${index}`),
idempotencyKey: id(`projection-start-${index}`),
fingerprint: id(`projection-fingerprint-${index}`),
now,
lease,
});
let job = await storage.jobs.enqueue(backgroundJob(channelId, lease, {
type: "research-answer",
targetKey: `research-answer:workspace:${workspaceId}:turn:${started.turn.id}:answer`,
availableAt: now,
now,
}));
await storage.research.linkJob({
channelId,
workspaceId,
turnId: started.turn.id,
role: "answer",
jobId: job.job.id,
now,
lease,
});
references.push({ id: workspaceId, jobId: job.job.id });
}

let commit = storage.collaboration.commit({
channelId,
lease,
expectedRevision: 0,
operationId: id("multi-projection-commit"),
epoch: "projection-epoch",
update: new Uint8Array([1]),
events: [],
now,
researchProjections: references.map(({ id }) => ({ id, action: "add" })),
});
let revisionBeforeTerminalization = Promise.resolve().then(async () => {
let state = await storage.collaboration.load(channelId, now);
await storage.jobs.cancel({
channelId,
jobId: references[0]!.jobId,
now,
lease,
});
return state?.channel.revision;
});

let [result, observedRevision] = await Promise.all([commit, revisionBeforeTerminalization]);
expect(observedRevision).toBe(result.revision);
expect(result.revision).toBe(1);
} finally {
await storage.close();
}
});
18 changes: 16 additions & 2 deletions apps/server/src/storage/memory/adapter.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@ import { documentSlug, documentSlugCandidate } from "../../channels/slug";
import { availableChannelTitle } from "../../channels/title";
import { MemoryBackgroundJobStore } from "./jobs";
import { MemoryResearchWorkspaceStore } from "./research";
import { researchProjectionAllowed, ResearchProjectionConflict } from "../model";

import type {
AddUserProject,
Expand Down Expand Up @@ -742,14 +743,27 @@ export class MemoryStorage implements StorageAdapter {
});
}

#commit(input: CommitChannel): Promise<CommitResult> {
async #commit(input: CommitChannel): Promise<CommitResult> {
this.#assertLease(input.lease);
let replay = this.#operations.get(input.channelId)?.get(input.operationId);
if (replay) return { ...replay, repeated: true };
for (let change of input.researchProjections ?? []) {
// Keep every projection check and the commit in one synchronous turn.
let detail = this.#research.getCurrent(input.channelId, change.id);
let initial = detail?.turns.find(turn => turn.kind === "initial");
let jobId = initial?.answerJobId ?? initial?.evidenceJobId;
let job = jobId ? this.#jobs.detail(input.channelId, jobId)?.job : undefined;
if (!researchProjectionAllowed(input.channelId, change, detail?.workspace, initial, job)) {
throw new ResearchProjectionConflict(change.id);
}
}
this.#assertLease(input.lease);
let found = this.#channels.get(input.channelId);
if (!found) throw missing(`channel ${input.channelId} does not exist`);
let operations = this.#operations.get(input.channelId) ?? new Map<string, Operation>();
this.#operations.set(input.channelId, operations);
let repeated = operations.get(input.operationId);
if (repeated) return Promise.resolve({ ...repeated, repeated: true });
if (repeated) return { ...repeated, repeated: true };
if (found.archivedAt && !input.allowArchived) {
throw conflict(`channel ${input.channelId} is archived`);
}
Expand Down
97 changes: 94 additions & 3 deletions apps/server/src/storage/memory/research.ts
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,7 @@ import type {
PublishInitialResearchReportResult,
ResearchJobRole,
ResearchMessage,
ResearchTerminalRecovery,
ResearchTurn,
ResearchWorkspace,
ResearchWorkspaceDetail,
Expand Down Expand Up @@ -189,6 +190,9 @@ export class MemoryResearchWorkspaceStore implements ResearchWorkspaceStore {
if (input.origin !== "inline" && input.origin !== "planner") {
throw conflict("research workspace start origin is invalid");
}
if (input.inlineReference && input.origin !== "planner") {
throw conflict("inline reference requires Planner origin");
}
if (
input.origin === "inline" && input.originMessageId !== undefined
|| input.origin === "planner" && !input.originMessageId
Expand All @@ -208,6 +212,9 @@ export class MemoryResearchWorkspaceStore implements ResearchWorkspaceStore {
`research workspace idempotency key ${input.idempotencyKey} was reused`,
);
}
if (input.inlineReference && !repeated.inlineReference) {
repeated.inlineReference = "pending";
}
return {
workspace: workspace(repeated),
turn: turn(repeatedTurn),
Expand Down Expand Up @@ -238,6 +245,7 @@ export class MemoryResearchWorkspaceStore implements ResearchWorkspaceStore {
confirmedQuery: input.question,
origin: input.origin,
originMessageId: input.originMessageId,
...(input.inlineReference ? { inlineReference: input.inlineReference } : {}),
createdBy: input.createdBy,
confirmedBy: input.createdBy,
revision: 0,
Expand Down Expand Up @@ -280,6 +288,70 @@ export class MemoryResearchWorkspaceStore implements ResearchWorkspaceStore {
};
});

readonly markReferencePlaced = (input: {
channelId: string;
workspaceId: string;
lease: Lease;
}): Promise<void> =>
this.#mutate(async () => {
this.#assertChannel(input.channelId, input.lease);
let found = this.#requireWorkspace(input.channelId, input.workspaceId);
if (!found.inlineReference) throw conflict("research reference placement is not required");
found.inlineReference = "placed";
});

readonly listReferenceRecovery = async (
limit: number,
afterId?: string,
channelId?: string,
): Promise<ResearchWorkspace[]> => {
let count = Math.min(100, Math.max(1, limit));
return [...this.#workspaces.values()]
.filter(value => !channelId || value.channelId === channelId)
.filter(value =>
value.inlineReference === "pending"
|| value.inlineReference === "placed"
&& !(this.#turns.get(value.id) ?? [])[0]?.evidenceJobId
)
.filter(value => !afterId || compareId(value.id, afterId) > 0)
.sort((left, right) => compareId(left.id, right.id))
.slice(0, count)
.map(workspace);
};

readonly listTerminalRecovery = async (
limit: number,
afterId?: string,
channelId?: string,
): Promise<ResearchTerminalRecovery[]> => {
let count = Math.min(100, Math.max(1, limit));
let candidates = [...this.#workspaces.values()]
.filter(value => value.origin === "planner" && value.inlineReference === "placed")
.filter(value => !channelId || value.channelId === channelId)
.filter(value => !afterId || compareId(value.id, afterId) > 0)
.sort((left, right) => compareId(left.id, right.id));
let page: ResearchTerminalRecovery[] = [];
for (let value of candidates) {
let initial = (this.#turns.get(value.id) ?? []).find(turn => turn.kind === "initial");
if (!initial) continue;
let answer = initial.answerJobId
? await this.#options.job(value.channelId, initial.answerJobId)
: undefined;
let evidence = initial.evidenceJobId
? await this.#options.job(value.channelId, initial.evidenceJobId)
: undefined;
let jobId: string | undefined;
if (
answer?.job.state === "failed"
|| answer?.job.state === "completed" && value.publishedChannelId
) jobId = answer.job.id;
else if (evidence?.job.state === "failed") jobId = evidence.job.id;
if (jobId) page.push({ id: value.id, channelId: value.channelId, jobId });
if (page.length === count) break;
}
return page;
};

readonly confirm = (
input: ConfirmResearchWorkspace,
): Promise<ConfirmResearchWorkspaceResult> =>
Expand Down Expand Up @@ -682,10 +754,14 @@ export class MemoryResearchWorkspaceStore implements ResearchWorkspaceStore {
readonly list = async (
channelId: string,
limit: number,
includePlanner = true,
): Promise<ResearchWorkspaceSummary[]> => {
let count = Math.min(100, Math.max(1, limit));
return [...this.#workspaces.values()]
.filter(value => value.channelId === channelId)
.filter(value =>
value.channelId === channelId
&& (includePlanner || value.origin !== "planner")
)
.sort((left, right) =>
right.updatedAt.getTime() - left.updatedAt.getTime() || compareId(left.id, right.id)
)
Expand All @@ -697,6 +773,7 @@ export class MemoryResearchWorkspaceStore implements ResearchWorkspaceStore {
repositoryId: string,
limit: number,
includeArchived = false,
includePlanner = true,
): Promise<ResearchWorkspaceRepositoryList> => {
let count = Math.min(RESEARCH_REPOSITORY_WORKSPACE_LIMIT, Math.max(1, limit));
let orderedChannels = this.#options.channels(repositoryId)
Expand All @@ -706,6 +783,7 @@ export class MemoryResearchWorkspaceStore implements ResearchWorkspaceStore {
);
let workspacesByChannel = new Map<string, ResearchWorkspace[]>();
for (let saved of this.#workspaces.values()) {
if (!includePlanner && saved.origin === "planner") continue;
let values = workspacesByChannel.get(saved.channelId) ?? [];
values.push(saved);
workspacesByChannel.set(saved.channelId, values);
Expand Down Expand Up @@ -744,17 +822,30 @@ export class MemoryResearchWorkspaceStore implements ResearchWorkspaceStore {
return { channels: groups, truncated };
};

readonly get = async (
getCurrent(
channelId: string,
workspaceId: string,
): Promise<ResearchWorkspaceDetail | undefined> => {
): ResearchWorkspaceDetail | undefined {
let found = this.#workspaces.get(workspaceId);
if (!found || found.channelId !== channelId) return undefined;
return {
workspace: workspace(found),
turns: (this.#turns.get(found.id) ?? []).map(turn),
messages: (this.#messages.get(found.id) ?? []).map(message),
};
}

readonly get = async (
channelId: string,
workspaceId: string,
): Promise<ResearchWorkspaceDetail | undefined> => this.getCurrent(channelId, workspaceId);

readonly findByIdempotencyKey = async (
channelId: string,
idempotencyKey: string,
): Promise<ResearchWorkspaceDetail | undefined> => {
let workspaceId = this.#idempotency.get(key(channelId, idempotencyKey));
return workspaceId ? this.get(channelId, workspaceId) : undefined;
};

readonly findTurnByJob = async (
Expand Down
Loading
Loading