mirror of
https://github.com/THU-MAIC/OpenMAIC.git
synced 2026-10-02 09:24:43 +08:00
feat(runtime): cut chat sessions over to RuntimeStore (#926)
* feat(runtime): persist chat sessions in RuntimeStore * fix(runtime): harden chat persistence ordering * fix(runtime): protect cross-tab chat snapshots * fix(runtime): compare structured chat values safely * fix(runtime): preserve chat backup and reset semantics * fix(runtime): restore imported chat snapshots * fix(runtime): isolate chat load failures * fix(runtime): make cache clearing lossless * fix(runtime): normalize empty chat saves * fix(runtime): serialize legacy chat migration * fix(chat): serialize backup restoration * test(chat): cover multi-stage restore locking * fix(chat): coordinate fallback migrations across tabs * fix(chat): harden cross-tab migration locking * fix(database): retire abandoned chat lease schema * fix(chat): preserve monotonic local edits * fix(chat): protect unseen runtime sessions * fix(chat): order lifecycle transitions monotonically * fix(chat): serialize runtime maintenance * fix(chat): order queued writes before maintenance * fix(chat): preserve cross-version lock compatibility * fix(chat): preserve streamed and backup data * fix(chat): isolate empty no-lock saves * fix(chat): isolate partial runtime snapshots * fix(runtime): close cross-store cutover races * fix(runtime): close remaining cutover races * fix(chat): retain caller-visible conflict baseline * fix(chat): preserve read-only legacy autosaves * fix(chat): avoid restore lock inversion * fix(runtime): enroll writes before maintenance * fix(runtime): retain maintenance ordering * fix(runtime): close restore and deletion races * fix(chat): order partition queues within lock epochs * fix(chat): preserve recovery and deletion intent * fix(runtime): recover interrupted chat restores * fix(backup): scope chat deduplication by stage * fix(backup): stage chats by runtime partition * fix(backup): clear staged chats with stages
This commit is contained in:
@@ -1,12 +1,16 @@
|
||||
'use client';
|
||||
|
||||
import { useState, useCallback, useRef, useEffect } from 'react';
|
||||
import type {
|
||||
ChatSession,
|
||||
SessionType,
|
||||
SessionStatus,
|
||||
ChatMessageMetadata,
|
||||
DirectorState,
|
||||
import {
|
||||
interruptActiveChatSessions,
|
||||
nextChatUpdatedAt,
|
||||
withChatSegmentReveal,
|
||||
withChatSessionStatus,
|
||||
type ChatSession,
|
||||
type SessionType,
|
||||
type SessionStatus,
|
||||
type ChatMessageMetadata,
|
||||
type DirectorState,
|
||||
} from '@/lib/types/chat';
|
||||
import type { DiscussionRequest } from '@/components/roundtable';
|
||||
import type { Action, SpotlightAction, DiscussionAction } from '@/lib/types/action';
|
||||
@@ -124,9 +128,7 @@ export function useChatSessions(options: UseChatSessionsOptions = {}) {
|
||||
const [sessions, setSessions] = useState<ChatSession[]>(() => {
|
||||
// Restore sessions from store (loaded from IndexedDB)
|
||||
const stored = useStageStore.getState().chats;
|
||||
return stored.map((s) =>
|
||||
s.status === 'active' ? { ...s, status: 'interrupted' as SessionStatus } : s,
|
||||
);
|
||||
return interruptActiveChatSessions(stored);
|
||||
});
|
||||
const [activeSessionId, setActiveSessionId] = useState<string | null>(null);
|
||||
const [expandedSessionIds, setExpandedSessionIds] = useState<Set<string>>(new Set());
|
||||
@@ -154,11 +156,7 @@ export function useChatSessions(options: UseChatSessionsOptions = {}) {
|
||||
stageIdRef.current = stageId;
|
||||
// Stage changed — reload sessions from store (already populated by loadFromStorage)
|
||||
const stored = useStageStore.getState().chats;
|
||||
setSessions(
|
||||
stored.map((s) =>
|
||||
s.status === 'active' ? { ...s, status: 'interrupted' as SessionStatus } : s,
|
||||
),
|
||||
);
|
||||
setSessions(interruptActiveChatSessions(stored));
|
||||
setActiveSessionId(null);
|
||||
setExpandedSessionIds(new Set());
|
||||
}, [stageId]);
|
||||
@@ -207,7 +205,7 @@ export function useChatSessions(options: UseChatSessionsOptions = {}) {
|
||||
? {
|
||||
...s,
|
||||
status: 'error' as SessionStatus,
|
||||
updatedAt: now,
|
||||
updatedAt: nextChatUpdatedAt(s, now),
|
||||
messages: [
|
||||
...s.messages,
|
||||
{
|
||||
@@ -289,7 +287,11 @@ export function useChatSessions(options: UseChatSessionsOptions = {}) {
|
||||
setSessions((prev) =>
|
||||
prev.map((s) =>
|
||||
s.id === sessionId
|
||||
? { ...s, messages: [...s.messages, newMsg], updatedAt: now }
|
||||
? {
|
||||
...s,
|
||||
messages: [...s.messages, newMsg],
|
||||
updatedAt: nextChatUpdatedAt(s, now),
|
||||
}
|
||||
: s,
|
||||
),
|
||||
);
|
||||
@@ -304,7 +306,9 @@ export function useChatSessions(options: UseChatSessionsOptions = {}) {
|
||||
const msgs = s.messages.filter(
|
||||
(m) => !(m.role === 'assistant' && m.parts.length === 0),
|
||||
);
|
||||
return msgs.length !== s.messages.length ? { ...s, messages: msgs } : s;
|
||||
return msgs.length !== s.messages.length
|
||||
? { ...s, messages: msgs, updatedAt: nextChatUpdatedAt(s) }
|
||||
: s;
|
||||
}),
|
||||
);
|
||||
},
|
||||
@@ -313,12 +317,12 @@ export function useChatSessions(options: UseChatSessionsOptions = {}) {
|
||||
messageId: string,
|
||||
partId: string,
|
||||
revealedText: string,
|
||||
_isComplete: boolean,
|
||||
isComplete: boolean,
|
||||
) {
|
||||
setSessions((prev) =>
|
||||
prev.map((s) => {
|
||||
if (s.id !== sessionId) return s;
|
||||
return {
|
||||
const revealed = {
|
||||
...s,
|
||||
messages: s.messages.map((m) => {
|
||||
if (m.id !== messageId) return m;
|
||||
@@ -344,6 +348,7 @@ export function useChatSessions(options: UseChatSessionsOptions = {}) {
|
||||
}),
|
||||
// Don't update updatedAt on every tick — avoids thrashing persistence sync
|
||||
};
|
||||
return withChatSegmentReveal(revealed, isComplete);
|
||||
}),
|
||||
);
|
||||
},
|
||||
@@ -367,7 +372,7 @@ export function useChatSessions(options: UseChatSessionsOptions = {}) {
|
||||
messages: s.messages.map((m) =>
|
||||
m.id === messageId ? { ...m, parts: [...m.parts, actionPart] } : m,
|
||||
),
|
||||
updatedAt: Date.now(),
|
||||
updatedAt: nextChatUpdatedAt(s),
|
||||
};
|
||||
}),
|
||||
);
|
||||
@@ -654,7 +659,11 @@ export function useChatSessions(options: UseChatSessionsOptions = {}) {
|
||||
setSessions((prev) =>
|
||||
prev.map((s) =>
|
||||
s.id === sessionId
|
||||
? { ...s, status: 'completed' as SessionStatus, updatedAt: Date.now() }
|
||||
? {
|
||||
...s,
|
||||
status: 'completed' as SessionStatus,
|
||||
updatedAt: nextChatUpdatedAt(s),
|
||||
}
|
||||
: s,
|
||||
),
|
||||
);
|
||||
@@ -775,7 +784,7 @@ export function useChatSessions(options: UseChatSessionsOptions = {}) {
|
||||
break;
|
||||
}
|
||||
}
|
||||
return { ...s, messages, status: 'completed' as SessionStatus };
|
||||
return { ...withChatSessionStatus(s, 'completed'), messages };
|
||||
}),
|
||||
);
|
||||
// Clear roundtable state via callbacks
|
||||
@@ -783,9 +792,7 @@ export function useChatSessions(options: UseChatSessionsOptions = {}) {
|
||||
onThinkingRef.current?.(null);
|
||||
} else {
|
||||
setSessions((prev) =>
|
||||
prev.map((s) =>
|
||||
s.id === sessionId ? { ...s, status: 'completed' as SessionStatus } : s,
|
||||
),
|
||||
prev.map((s) => (s.id === sessionId ? withChatSessionStatus(s, 'completed') : s)),
|
||||
);
|
||||
}
|
||||
|
||||
@@ -876,7 +883,7 @@ export function useChatSessions(options: UseChatSessionsOptions = {}) {
|
||||
}
|
||||
}
|
||||
// Keep status 'active' — session continues when user speaks
|
||||
return { ...s, messages, updatedAt: Date.now() };
|
||||
return { ...s, messages, updatedAt: nextChatUpdatedAt(s) };
|
||||
}),
|
||||
);
|
||||
// Note: Do NOT call onLiveSpeech/onThinking here.
|
||||
@@ -1021,7 +1028,7 @@ export function useChatSessions(options: UseChatSessionsOptions = {}) {
|
||||
text: (textPart.text || '') + '...',
|
||||
} as UIMessage<ChatMessageMetadata>['parts'][number];
|
||||
messages[i] = { ...messages[i], parts };
|
||||
return { ...s, messages, updatedAt: Date.now() };
|
||||
return { ...s, messages, updatedAt: nextChatUpdatedAt(s) };
|
||||
}
|
||||
}
|
||||
break;
|
||||
@@ -1105,7 +1112,7 @@ export function useChatSessions(options: UseChatSessionsOptions = {}) {
|
||||
...s,
|
||||
messages: [...s.messages, userMessage],
|
||||
status: 'active' as SessionStatus,
|
||||
updatedAt: now,
|
||||
updatedAt: nextChatUpdatedAt(s, now),
|
||||
}
|
||||
: s,
|
||||
);
|
||||
@@ -1373,9 +1380,7 @@ export function useChatSessions(options: UseChatSessionsOptions = {}) {
|
||||
// Actions won't be re-appended because lastActionIndex already covers them.
|
||||
if (existing.status === 'completed') {
|
||||
setSessions((prev) =>
|
||||
prev.map((s) =>
|
||||
s.id === existing.id ? { ...s, status: 'active' as SessionStatus } : s,
|
||||
),
|
||||
prev.map((s) => (s.id === existing.id ? withChatSessionStatus(s, 'active') : s)),
|
||||
);
|
||||
// Restore lecture tracking refs (cleared by endSession)
|
||||
const messageId = existing.messages[0]?.id;
|
||||
@@ -1461,7 +1466,9 @@ export function useChatSessions(options: UseChatSessionsOptions = {}) {
|
||||
// Update lastActionIndex in session
|
||||
setSessions((prev) =>
|
||||
prev.map((s) =>
|
||||
s.id === sessionId ? { ...s, lastActionIndex: actionIndex, updatedAt: Date.now() } : s,
|
||||
s.id === sessionId
|
||||
? { ...s, lastActionIndex: actionIndex, updatedAt: nextChatUpdatedAt(s) }
|
||||
: s,
|
||||
),
|
||||
);
|
||||
|
||||
|
||||
@@ -1,7 +1,13 @@
|
||||
import { restoreAgentSelection } from '@/lib/orchestration/registry/agent-selection';
|
||||
import { useAgentRegistry } from '@/lib/orchestration/registry/store';
|
||||
import { getActionsForRole } from '@/lib/orchestration/registry/types';
|
||||
import { applyHydratedClassroomFallbackScenes } from '@/lib/classroom/pbl-fallback-hydration';
|
||||
import {
|
||||
applyHydratedClassroomFallbackScenes,
|
||||
hydrateClassroomFallbackChats,
|
||||
type ApplyHydratedClassroomFallbackScenesArgs,
|
||||
} from '@/lib/classroom/pbl-fallback-hydration';
|
||||
import type { ChatStorageSnapshot } from '@/lib/utils/chat-storage';
|
||||
import type { ChatSession } from '@/lib/types/chat';
|
||||
import type { TTSProviderId } from '@/lib/audio/types';
|
||||
import type { VoiceDesign } from '@/lib/audio/voice-design';
|
||||
import { useMediaGenerationStore, type MediaTask } from '@/lib/store/media-generation';
|
||||
@@ -175,14 +181,19 @@ export async function fetchClassroomFromApi(classroomId: string): Promise<Classr
|
||||
export function applyClassroomStageAndScenes(
|
||||
stage: Stage,
|
||||
scenes: readonly Scene[],
|
||||
options: { persist?: boolean } = {},
|
||||
options: {
|
||||
persist?: boolean;
|
||||
chats?: ChatSession[];
|
||||
chatSnapshot?: ChatStorageSnapshot;
|
||||
} = {},
|
||||
): void {
|
||||
const nextScenes = [...scenes];
|
||||
useStageStore.setState((state) => ({
|
||||
stage,
|
||||
scenes: nextScenes,
|
||||
currentSceneId: nextScenes[0]?.id ?? null,
|
||||
chats: [],
|
||||
chats: options.chats ?? [],
|
||||
chatSnapshot: options.chatSnapshot ?? { sessions: [], restoreMarker: null },
|
||||
generationComplete: false,
|
||||
generationEpoch: state.generationEpoch + 1,
|
||||
mode: 'playback',
|
||||
@@ -346,7 +357,11 @@ export function applyGeneratedAgentRecordsToRegistry(
|
||||
}
|
||||
|
||||
export const defaultClassroomLoadDeps = {
|
||||
applyFallbackScenes: applyHydratedClassroomFallbackScenes,
|
||||
applyFallbackScenes: (args: ApplyHydratedClassroomFallbackScenesArgs) =>
|
||||
applyHydratedClassroomFallbackScenes({
|
||||
...args,
|
||||
hydrateChats: hydrateClassroomFallbackChats,
|
||||
}),
|
||||
fetchClassroom: fetchClassroomFromApi,
|
||||
loadRestoredMediaTasks: loadRestoredMediaTasksFromDB,
|
||||
applyRestoredMediaTasks,
|
||||
|
||||
@@ -4,7 +4,13 @@ import {
|
||||
type HydratePBLProjectArgs,
|
||||
} from '@/lib/pbl/v2/runtime/hydration';
|
||||
import { isCurrentStageSceneLoadToken, type StageSceneLoadToken } from '@/lib/store/stage';
|
||||
import type { ChatSession } from '@/lib/types/chat';
|
||||
import type { Scene, Stage } from '@/lib/types/stage';
|
||||
import {
|
||||
loadChatSessions,
|
||||
type ChatStorageReadOptions,
|
||||
type ChatStorageSnapshot,
|
||||
} from '@/lib/utils/chat-storage';
|
||||
|
||||
export async function hydrateClassroomFallbackScenes(
|
||||
stageId: string,
|
||||
@@ -14,13 +20,39 @@ export async function hydrateClassroomFallbackScenes(
|
||||
return hydratePBLScenesFromRuntime(stageId, scenes.map(migrateScene), options);
|
||||
}
|
||||
|
||||
export interface ClassroomFallbackChatState {
|
||||
chats: ChatSession[];
|
||||
chatSnapshot: ChatStorageSnapshot;
|
||||
}
|
||||
|
||||
export async function hydrateClassroomFallbackChats(
|
||||
stageId: string,
|
||||
options: ChatStorageReadOptions = {},
|
||||
): Promise<ClassroomFallbackChatState> {
|
||||
let chatSnapshot: ChatStorageSnapshot = { sessions: [], restoreMarker: undefined };
|
||||
try {
|
||||
const chats = await loadChatSessions(stageId, {
|
||||
...options,
|
||||
onSnapshot: (snapshot) => {
|
||||
chatSnapshot = snapshot;
|
||||
options.onSnapshot?.(snapshot);
|
||||
},
|
||||
});
|
||||
return { chats, chatSnapshot };
|
||||
} catch (error) {
|
||||
console.warn(`Failed to hydrate runtime chats for server fallback stage ${stageId}:`, error);
|
||||
return { chats: [], chatSnapshot };
|
||||
}
|
||||
}
|
||||
|
||||
export interface ApplyHydratedClassroomFallbackScenesArgs {
|
||||
loadToken: StageSceneLoadToken;
|
||||
isCurrent?: () => boolean;
|
||||
stage: Stage;
|
||||
scenes: readonly Scene[];
|
||||
hydrateScenes?: (stageId: string, scenes: readonly Scene[]) => Promise<Scene[]>;
|
||||
applyStageAndScenes: (stage: Stage, scenes: Scene[]) => void;
|
||||
hydrateChats?: (stageId: string) => Promise<ClassroomFallbackChatState>;
|
||||
applyStageAndScenes: (stage: Stage, scenes: Scene[], options: ClassroomFallbackChatState) => void;
|
||||
}
|
||||
|
||||
export async function applyHydratedClassroomFallbackScenes({
|
||||
@@ -29,12 +61,19 @@ export async function applyHydratedClassroomFallbackScenes({
|
||||
stage,
|
||||
scenes,
|
||||
hydrateScenes = hydrateClassroomFallbackScenes,
|
||||
hydrateChats = async () => ({
|
||||
chats: [],
|
||||
chatSnapshot: { sessions: [], restoreMarker: null },
|
||||
}),
|
||||
applyStageAndScenes,
|
||||
}: ApplyHydratedClassroomFallbackScenesArgs): Promise<boolean> {
|
||||
const hydrated = await hydrateScenes(stage.id, scenes);
|
||||
const [hydrated, chatState] = await Promise.all([
|
||||
hydrateScenes(stage.id, scenes),
|
||||
hydrateChats(stage.id),
|
||||
]);
|
||||
if (!isCurrent() || !isCurrentStageSceneLoadToken(loadToken)) {
|
||||
return false;
|
||||
}
|
||||
applyStageAndScenes(stage, hydrated);
|
||||
applyStageAndScenes(stage, hydrated, chatState);
|
||||
return true;
|
||||
}
|
||||
|
||||
+66
-29
@@ -22,12 +22,13 @@ import { BrowserKVStore, type KVStore, type RuntimeStore } from '@openmaic/stora
|
||||
|
||||
import { getLearnerKey } from '@/lib/runtime/learner-key';
|
||||
import { getRuntimeStore } from '@/lib/runtime/store';
|
||||
import { withRuntimeStorageSharedLock } from '@/lib/utils/chat-storage-lock';
|
||||
import type { PBLEngagementEvent, PBLProjectV2, PBLRuntimeEvent } from '@/lib/pbl/v2/types';
|
||||
import { enrichPBLRuntimeEvent, pblEngagementRecordPayload } from './record-payloads';
|
||||
|
||||
const PBL_DRAIN_TIMEOUT_MS = 10_000;
|
||||
const PBL_DRAIN_CHAIN_HARD_CAP_MS = PBL_DRAIN_TIMEOUT_MS * 2;
|
||||
const PBL_HYDRATION_DRAIN_BARRIER_TIMEOUT_MS = PBL_DRAIN_CHAIN_HARD_CAP_MS;
|
||||
export const PBL_HYDRATION_DRAIN_BARRIER_TIMEOUT_MS = PBL_DRAIN_CHAIN_HARD_CAP_MS;
|
||||
const WATERMARK_SCOPE = 'device';
|
||||
|
||||
let defaultKv: KVStore | undefined;
|
||||
@@ -74,10 +75,10 @@ function isDeadlineExpired(deadline: PBLDrainDeadline): boolean {
|
||||
}
|
||||
|
||||
/** Reject after `ms`, clearing the timer once the raced promise settles. */
|
||||
async function withTimeout(work: Promise<void>, ms: number): Promise<void> {
|
||||
async function withTimeout<T>(work: Promise<T>, ms: number): Promise<T> {
|
||||
let timer: ReturnType<typeof setTimeout> | undefined;
|
||||
try {
|
||||
await Promise.race([
|
||||
return await Promise.race([
|
||||
work,
|
||||
new Promise<never>((_, reject) => {
|
||||
timer = setTimeout(() => reject(new Error(`timed out after ${ms}ms`)), ms);
|
||||
@@ -325,10 +326,7 @@ async function waitForActiveDrainWork(key: string): Promise<void> {
|
||||
}
|
||||
}
|
||||
|
||||
async function drainProjectRuntimeSerialized(
|
||||
args: DrainProjectRuntimeArgs,
|
||||
waitForActualCompletion = false,
|
||||
): Promise<void> {
|
||||
async function drainProjectRuntimeSerialized(args: DrainProjectRuntimeArgs): Promise<void> {
|
||||
const kv = args.kv ?? getDefaultKv();
|
||||
const learnerKey = args.learnerKey ?? (await getLearnerKey(kv));
|
||||
const store = args.store ?? getRuntimeStore();
|
||||
@@ -340,15 +338,7 @@ async function drainProjectRuntimeSerialized(
|
||||
// needs to re-read the watermark and make progress from the durable point.
|
||||
})
|
||||
.then(async () => {
|
||||
if (waitForActualCompletion) {
|
||||
await withTimeout(
|
||||
waitForActiveDrainWork(inFlightKey),
|
||||
PBL_HYDRATION_DRAIN_BARRIER_TIMEOUT_MS,
|
||||
);
|
||||
}
|
||||
const deadline = createDrainDeadline(
|
||||
waitForActualCompletion ? Number.POSITIVE_INFINITY : PBL_DRAIN_TIMEOUT_MS,
|
||||
);
|
||||
const deadline = createDrainDeadline(PBL_DRAIN_TIMEOUT_MS);
|
||||
const drainWork = drainProjectRuntimeWork(
|
||||
{
|
||||
...args,
|
||||
@@ -359,9 +349,6 @@ async function drainProjectRuntimeSerialized(
|
||||
deadline,
|
||||
);
|
||||
trackActiveDrainWork(inFlightKey, drainWork);
|
||||
if (waitForActualCompletion) {
|
||||
return withTimeout(drainWork, PBL_HYDRATION_DRAIN_BARRIER_TIMEOUT_MS);
|
||||
}
|
||||
// Safety invariant for releasing a stuck chain link: drain work checks
|
||||
// the cooperative deadline before each append and before watermark
|
||||
// writes. After that deadline, an overdue append may still complete, but
|
||||
@@ -381,21 +368,71 @@ async function drainProjectRuntimeSerialized(
|
||||
}
|
||||
|
||||
/**
|
||||
* Drain all currently visible events behind a bounded completion barrier.
|
||||
* Hydration and document persistence must either observe the completed drain
|
||||
* or abort; neither may fold or write a snapshot while an append is pending.
|
||||
* Drain all currently visible events behind a bounded completion barrier, then
|
||||
* keep the runtime-wide shared lock through the caller's fold/snapshot work.
|
||||
* The caller enrolls in the shared maintenance epoch before waiting, so work
|
||||
* already requested cannot resume after a later destructive operation.
|
||||
*/
|
||||
export async function drainProjectRuntimeFully(args: DrainProjectRuntimeArgs): Promise<void> {
|
||||
const work = drainProjectRuntimeSerialized(args, true);
|
||||
// The budget covers queueing behind prior saves as well as this drain's own
|
||||
// work. A late rejection is observed here after the caller has fallen back.
|
||||
work.catch(() => {});
|
||||
await withTimeout(work, PBL_HYDRATION_DRAIN_BARRIER_TIMEOUT_MS);
|
||||
export async function withDrainedProjectRuntime<T>(
|
||||
args: DrainProjectRuntimeArgs,
|
||||
work: () => Promise<T>,
|
||||
globalLockHeld = false,
|
||||
): Promise<T> {
|
||||
const run = async (): Promise<T> => {
|
||||
const kv = args.kv ?? getDefaultKv();
|
||||
const learnerKey = args.learnerKey ?? (await getLearnerKey(kv));
|
||||
const store = args.store ?? getRuntimeStore();
|
||||
const inFlightKey = `${args.stageId}:${args.sceneId}:${learnerKey}`;
|
||||
const previous = inFlightPblDrains.get(inFlightKey) ?? Promise.resolve();
|
||||
const deadline = createDrainDeadline(PBL_HYDRATION_DRAIN_BARRIER_TIMEOUT_MS);
|
||||
const serialized = previous
|
||||
.catch(() => {
|
||||
// The earlier caller already observed its failure. Re-read the durable
|
||||
// watermark after any actual late append work settles.
|
||||
})
|
||||
.then(async () => {
|
||||
await withTimeout(
|
||||
waitForActiveDrainWork(inFlightKey),
|
||||
PBL_HYDRATION_DRAIN_BARRIER_TIMEOUT_MS,
|
||||
);
|
||||
await drainProjectRuntimeWork(
|
||||
{
|
||||
...args,
|
||||
store,
|
||||
kv,
|
||||
learnerKey,
|
||||
},
|
||||
deadline,
|
||||
);
|
||||
if (isDeadlineExpired(deadline)) {
|
||||
throw new Error(`timed out after ${PBL_HYDRATION_DRAIN_BARRIER_TIMEOUT_MS}ms`);
|
||||
}
|
||||
return work();
|
||||
});
|
||||
const chain = serialized.then(() => undefined);
|
||||
inFlightPblDrains.set(inFlightKey, chain);
|
||||
void chain
|
||||
.finally(() => {
|
||||
if (inFlightPblDrains.get(inFlightKey) === chain) {
|
||||
inFlightPblDrains.delete(inFlightKey);
|
||||
}
|
||||
})
|
||||
.catch(() => {});
|
||||
return serialized;
|
||||
};
|
||||
const operation = globalLockHeld ? run() : withRuntimeStorageSharedLock(run);
|
||||
operation.catch(() => {});
|
||||
return globalLockHeld
|
||||
? operation
|
||||
: withTimeout(operation, PBL_HYDRATION_DRAIN_BARRIER_TIMEOUT_MS);
|
||||
}
|
||||
|
||||
export async function drainProjectRuntime(args: DrainProjectRuntimeArgs): Promise<void> {
|
||||
try {
|
||||
const work = drainProjectRuntimeSerialized(args);
|
||||
// Enroll before any async context lookup or same-key queue wait. Otherwise
|
||||
// maintenance can overtake an already-requested drain and stale work can
|
||||
// recreate runtime rows after deletion.
|
||||
const work = withRuntimeStorageSharedLock(() => drainProjectRuntimeSerialized(args));
|
||||
// A rejection landing after the timeout already won the race would have no
|
||||
// listener left. Swallow that branch; the await below still reports it if it
|
||||
// lands before the timeout.
|
||||
|
||||
+135
-118
@@ -4,9 +4,14 @@ import { isEqual } from 'lodash';
|
||||
|
||||
import { getLearnerKey } from '@/lib/runtime/learner-key';
|
||||
import { getRuntimeStore } from '@/lib/runtime/store';
|
||||
import { withRuntimeStorageSharedLockUntilSettled } from '@/lib/utils/chat-storage-lock';
|
||||
import type { Scene } from '@/lib/types/stage';
|
||||
import type { PBLProjectV2 } from '@/lib/pbl/v2/types';
|
||||
import { drainProjectRuntimeFully, ensurePBLRuntimeSession } from './drain';
|
||||
import {
|
||||
ensurePBLRuntimeSession,
|
||||
PBL_HYDRATION_DRAIN_BARRIER_TIMEOUT_MS,
|
||||
withDrainedProjectRuntime,
|
||||
} from './drain';
|
||||
import { foldPBLRuntime, type PBLFoldDiagnostics } from './fold';
|
||||
import {
|
||||
applyLearnerState,
|
||||
@@ -173,138 +178,150 @@ function hasWriteCutoverSnapshot(records: readonly RuntimeRecord[]): boolean {
|
||||
}
|
||||
|
||||
export async function synchronizePBLProjectRuntime(args: HydratePBLProjectArgs): Promise<void> {
|
||||
const kv = args.kv ?? getDefaultKv();
|
||||
const learnerKey = args.learnerKey ?? (await getLearnerKey(kv));
|
||||
const store = args.store ?? getRuntimeStore();
|
||||
const transactionKey = `${args.stageId}:${args.sceneId}:${learnerKey}`;
|
||||
await withRuntimeStorageSharedLockUntilSettled(async () => {
|
||||
const kv = args.kv ?? getDefaultKv();
|
||||
const learnerKey = args.learnerKey ?? (await getLearnerKey(kv));
|
||||
const store = args.store ?? getRuntimeStore();
|
||||
const transactionKey = `${args.stageId}:${args.sceneId}:${learnerKey}`;
|
||||
|
||||
await withPBLRuntimeTransaction(store, transactionKey, async () => {
|
||||
await drainProjectRuntimeFully({
|
||||
stageId: args.stageId,
|
||||
sceneId: args.sceneId,
|
||||
project: args.project,
|
||||
store,
|
||||
kv,
|
||||
learnerKey,
|
||||
});
|
||||
await withPBLRuntimeTransaction(store, transactionKey, () =>
|
||||
withDrainedProjectRuntime(
|
||||
{
|
||||
stageId: args.stageId,
|
||||
sceneId: args.sceneId,
|
||||
project: args.project,
|
||||
store,
|
||||
kv,
|
||||
learnerKey,
|
||||
},
|
||||
async () => {
|
||||
const records = await listPBLRecords({
|
||||
store,
|
||||
stageId: args.stageId,
|
||||
sceneId: args.sceneId,
|
||||
learnerKey,
|
||||
});
|
||||
const learnerState = extractLearnerState(args.project);
|
||||
const folded = foldPBLRuntime({
|
||||
designTemplate: stripToDesignTemplate(args.project),
|
||||
records,
|
||||
});
|
||||
const runtimeIsCurrent =
|
||||
isEqual(folded.learnerState, learnerState) && folded.diagnostics.gaps.length === 0;
|
||||
const cutoverStarted = hasWriteCutoverSnapshot(records);
|
||||
if (runtimeIsCurrent && cutoverStarted) return;
|
||||
|
||||
const records = await listPBLRecords({
|
||||
store,
|
||||
stageId: args.stageId,
|
||||
sceneId: args.sceneId,
|
||||
learnerKey,
|
||||
});
|
||||
const learnerState = extractLearnerState(args.project);
|
||||
const folded = foldPBLRuntime({
|
||||
designTemplate: stripToDesignTemplate(args.project),
|
||||
records,
|
||||
});
|
||||
const runtimeIsCurrent =
|
||||
isEqual(folded.learnerState, learnerState) && folded.diagnostics.gaps.length === 0;
|
||||
const cutoverStarted = hasWriteCutoverSnapshot(records);
|
||||
if (runtimeIsCurrent && cutoverStarted) return;
|
||||
|
||||
await appendPBLRuntimeSnapshotIfChanged({
|
||||
store,
|
||||
stageId: args.stageId,
|
||||
sceneId: args.sceneId,
|
||||
learnerKey,
|
||||
project: args.project,
|
||||
learnerState,
|
||||
records,
|
||||
reason: cutoverStarted ? 'self_heal' : 'write_cutover',
|
||||
});
|
||||
});
|
||||
await appendPBLRuntimeSnapshotIfChanged({
|
||||
store,
|
||||
stageId: args.stageId,
|
||||
sceneId: args.sceneId,
|
||||
learnerKey,
|
||||
project: args.project,
|
||||
learnerState,
|
||||
records,
|
||||
reason: cutoverStarted ? 'self_heal' : 'write_cutover',
|
||||
});
|
||||
},
|
||||
true,
|
||||
),
|
||||
);
|
||||
}, PBL_HYDRATION_DRAIN_BARRIER_TIMEOUT_MS);
|
||||
}
|
||||
|
||||
export async function hydratePBLProjectFromRuntime(
|
||||
args: HydratePBLProjectArgs,
|
||||
): Promise<HydratePBLProjectResult> {
|
||||
const kv = args.kv ?? getDefaultKv();
|
||||
const learnerKey = args.learnerKey ?? (await getLearnerKey(kv));
|
||||
const store = args.store ?? getRuntimeStore();
|
||||
const transactionKey = `${args.stageId}:${args.sceneId}:${learnerKey}`;
|
||||
return withRuntimeStorageSharedLockUntilSettled(async () => {
|
||||
const kv = args.kv ?? getDefaultKv();
|
||||
const learnerKey = args.learnerKey ?? (await getLearnerKey(kv));
|
||||
const store = args.store ?? getRuntimeStore();
|
||||
const transactionKey = `${args.stageId}:${args.sceneId}:${learnerKey}`;
|
||||
|
||||
return withPBLRuntimeTransaction(store, transactionKey, async () => {
|
||||
await drainProjectRuntimeFully({
|
||||
stageId: args.stageId,
|
||||
sceneId: args.sceneId,
|
||||
project: args.project,
|
||||
store,
|
||||
kv,
|
||||
learnerKey,
|
||||
});
|
||||
return withPBLRuntimeTransaction(store, transactionKey, () =>
|
||||
withDrainedProjectRuntime(
|
||||
{
|
||||
stageId: args.stageId,
|
||||
sceneId: args.sceneId,
|
||||
project: args.project,
|
||||
store,
|
||||
kv,
|
||||
learnerKey,
|
||||
},
|
||||
async () => {
|
||||
const records = await listPBLRecords({
|
||||
store,
|
||||
stageId: args.stageId,
|
||||
sceneId: args.sceneId,
|
||||
learnerKey,
|
||||
});
|
||||
const designTemplate = stripToDesignTemplate(args.project);
|
||||
const folded = foldPBLRuntime({ designTemplate, records });
|
||||
const documentState = extractLearnerState(args.project);
|
||||
const stateMatchesDocument = isEqual(folded.learnerState, documentState);
|
||||
const matchesDocument = stateMatchesDocument && folded.diagnostics.gaps.length === 0;
|
||||
|
||||
const records = await listPBLRecords({
|
||||
store,
|
||||
stageId: args.stageId,
|
||||
sceneId: args.sceneId,
|
||||
learnerKey,
|
||||
});
|
||||
const designTemplate = stripToDesignTemplate(args.project);
|
||||
const folded = foldPBLRuntime({ designTemplate, records });
|
||||
const documentState = extractLearnerState(args.project);
|
||||
const stateMatchesDocument = isEqual(folded.learnerState, documentState);
|
||||
const matchesDocument = stateMatchesDocument && folded.diagnostics.gaps.length === 0;
|
||||
if (matchesDocument) {
|
||||
return {
|
||||
project: preserveDocumentTransients(
|
||||
applyLearnerState(designTemplate, folded.learnerState),
|
||||
args.project,
|
||||
),
|
||||
source: 'fold' as const,
|
||||
diagnostics: folded.diagnostics,
|
||||
diff: [],
|
||||
selfHealed: false,
|
||||
};
|
||||
}
|
||||
|
||||
if (matchesDocument) {
|
||||
return {
|
||||
project: preserveDocumentTransients(
|
||||
applyLearnerState(designTemplate, folded.learnerState),
|
||||
args.project,
|
||||
),
|
||||
source: 'fold' as const,
|
||||
diagnostics: folded.diagnostics,
|
||||
diff: [],
|
||||
selfHealed: false,
|
||||
};
|
||||
}
|
||||
const diff = diffLearnerState(folded.learnerState, documentState);
|
||||
const hasDocumentLearnerState = documentContainsLearnerState(args.project);
|
||||
const cutoverStarted = hasWriteCutoverSnapshot(records);
|
||||
|
||||
const diff = diffLearnerState(folded.learnerState, documentState);
|
||||
const hasDocumentLearnerState = documentContainsLearnerState(args.project);
|
||||
const cutoverStarted = hasWriteCutoverSnapshot(records);
|
||||
if (hasDocumentLearnerState && !cutoverStarted && process.env.NODE_ENV !== 'production') {
|
||||
console.warn('[PBL runtime] document state remained authoritative during hydration', {
|
||||
stageId: args.stageId,
|
||||
sceneId: args.sceneId,
|
||||
diff,
|
||||
gaps: folded.diagnostics.gaps.slice(0, 8),
|
||||
});
|
||||
}
|
||||
|
||||
if (hasDocumentLearnerState && !cutoverStarted && process.env.NODE_ENV !== 'production') {
|
||||
console.warn('[PBL runtime] document state remained authoritative during hydration', {
|
||||
stageId: args.stageId,
|
||||
sceneId: args.sceneId,
|
||||
diff,
|
||||
gaps: folded.diagnostics.gaps.slice(0, 8),
|
||||
});
|
||||
}
|
||||
if (hasDocumentLearnerState && !cutoverStarted) {
|
||||
const selfHealed = await appendPBLRuntimeSnapshotIfChanged({
|
||||
store,
|
||||
stageId: args.stageId,
|
||||
sceneId: args.sceneId,
|
||||
learnerKey,
|
||||
project: args.project,
|
||||
learnerState: documentState,
|
||||
records,
|
||||
reason: records.length === 0 ? 'backfill' : 'self_heal',
|
||||
});
|
||||
|
||||
if (hasDocumentLearnerState && !cutoverStarted) {
|
||||
const selfHealed = await appendPBLRuntimeSnapshotIfChanged({
|
||||
store,
|
||||
stageId: args.stageId,
|
||||
sceneId: args.sceneId,
|
||||
learnerKey,
|
||||
project: args.project,
|
||||
learnerState: documentState,
|
||||
records,
|
||||
reason: records.length === 0 ? 'backfill' : 'self_heal',
|
||||
});
|
||||
return {
|
||||
project: args.project,
|
||||
source: 'document' as const,
|
||||
diagnostics: folded.diagnostics,
|
||||
diff,
|
||||
selfHealed,
|
||||
};
|
||||
}
|
||||
|
||||
return {
|
||||
project: args.project,
|
||||
source: 'document' as const,
|
||||
diagnostics: folded.diagnostics,
|
||||
diff,
|
||||
selfHealed,
|
||||
};
|
||||
}
|
||||
|
||||
return {
|
||||
project: preserveDocumentTransients(
|
||||
applyLearnerState(designTemplate, folded.learnerState),
|
||||
args.project,
|
||||
return {
|
||||
project: preserveDocumentTransients(
|
||||
applyLearnerState(designTemplate, folded.learnerState),
|
||||
args.project,
|
||||
),
|
||||
source: 'fold' as const,
|
||||
diagnostics: folded.diagnostics,
|
||||
diff,
|
||||
selfHealed: false,
|
||||
};
|
||||
},
|
||||
true,
|
||||
),
|
||||
source: 'fold' as const,
|
||||
diagnostics: folded.diagnostics,
|
||||
diff,
|
||||
selfHealed: false,
|
||||
};
|
||||
});
|
||||
);
|
||||
}, PBL_HYDRATION_DRAIN_BARRIER_TIMEOUT_MS);
|
||||
}
|
||||
|
||||
export async function hydratePBLScenesFromRuntime(
|
||||
|
||||
+33
-21
@@ -1,8 +1,8 @@
|
||||
/**
|
||||
* Lazy app-wide RuntimeStore singleton (#869). One `maic-runtime` IndexedDB
|
||||
* per origin, shared by every runtime kind (pbl, chat, quizAttempt, playback)
|
||||
* as they migrate onto the runtime layer. Nothing reads or writes it yet
|
||||
* except the stage-deletion cascade; Part C2 adds the first real writer.
|
||||
* as they migrate onto the runtime layer. PBL events and chat sessions use it
|
||||
* today; the stage-deletion cascade clears every kind together.
|
||||
*
|
||||
* Client-only: the store lazily opens IndexedDB. Server code must not import
|
||||
* this module without injecting its own `RuntimeStore`.
|
||||
@@ -58,29 +58,41 @@ async function withTimeout(work: Promise<void>, ms: number): Promise<void> {
|
||||
* (`maic-runtime`), so a broken or hung runtime DB must not brick stage
|
||||
* deletion in the main app DB — the cascade is bounded by a timeout, and any
|
||||
* failure warns and moves on. A failed or timed-out cascade leaves orphaned
|
||||
* runtime rows, which are inert today (nothing reads them yet); a startup
|
||||
* sweep is deliberately deferred to Part C2, when the store gains real
|
||||
* readers.
|
||||
* runtime rows; they are not reachable through normal navigation once the
|
||||
* stage is gone, and a future startup sweep can reclaim them.
|
||||
*/
|
||||
export async function deleteStageRuntimeSafely(
|
||||
stageId: string,
|
||||
runtimeStore?: RuntimeStore,
|
||||
): Promise<void> {
|
||||
try {
|
||||
// Probe + cascade share the try/catch and the timeout envelope: a hanging
|
||||
// `databases()` must not brick deletion any more than a hanging store.
|
||||
const work = (async () => {
|
||||
if (!(await runtimeDbExists())) return; // nothing to clean
|
||||
const cascade = (runtimeStore ?? getRuntimeStore()).deleteStageRuntime(stageId);
|
||||
// A rejection landing after the timeout already won the race would have
|
||||
// no listener left — swallow that branch so it cannot surface as an
|
||||
// unhandled rejection (the await below still reports it if it lands in
|
||||
// time).
|
||||
cascade.catch(() => {});
|
||||
await cascade;
|
||||
})();
|
||||
await withTimeout(work, STAGE_RUNTIME_DELETE_TIMEOUT_MS);
|
||||
} catch (error) {
|
||||
await beginStageRuntimeDeletionSafely(stageId, runtimeStore).completion;
|
||||
}
|
||||
|
||||
export interface StageRuntimeDeletion {
|
||||
/** Bounded, fail-soft caller-visible completion. */
|
||||
completion: Promise<void>;
|
||||
/** Fail-soft actual settlement, used to retain destructive maintenance locks. */
|
||||
settlement: Promise<void>;
|
||||
}
|
||||
|
||||
/** Start one bounded deletion while keeping a handle to its real settlement. */
|
||||
export function beginStageRuntimeDeletionSafely(
|
||||
stageId: string,
|
||||
runtimeStore?: RuntimeStore,
|
||||
): StageRuntimeDeletion {
|
||||
// Probe + cascade share the same underlying work: a hanging databases()
|
||||
// probe is just as important to retain behind maintenance as a hanging delete.
|
||||
const work = (async () => {
|
||||
if (!(await runtimeDbExists())) return;
|
||||
await (runtimeStore ?? getRuntimeStore()).deleteStageRuntime(stageId);
|
||||
})();
|
||||
let reported = false;
|
||||
const report = (error: unknown): void => {
|
||||
if (reported) return;
|
||||
reported = true;
|
||||
console.warn(`Failed to delete runtime data for stage ${stageId}:`, error);
|
||||
}
|
||||
};
|
||||
const settlement = work.catch(report);
|
||||
const completion = withTimeout(work, STAGE_RUNTIME_DELETE_TIMEOUT_MS).catch(report);
|
||||
return { completion, settlement };
|
||||
}
|
||||
|
||||
+20
-1
@@ -17,6 +17,7 @@ import { useCanvasStore } from '@/lib/store/canvas';
|
||||
import { migrateScene } from '@/lib/edit/slide-schema';
|
||||
import { preparePBLScenesForDocumentPersistence } from '@/lib/pbl/v2/runtime/document-persistence';
|
||||
import { hydratePBLScenesFromRuntime } from '@/lib/pbl/v2/runtime/hydration';
|
||||
import type { ChatStorageSnapshot } from '@/lib/utils/chat-storage';
|
||||
|
||||
const log = createLogger('StageStore');
|
||||
|
||||
@@ -87,6 +88,7 @@ interface StageState {
|
||||
|
||||
// Chats
|
||||
chats: ChatSession[];
|
||||
chatSnapshot: ChatStorageSnapshot;
|
||||
|
||||
// Mode
|
||||
mode: StageMode;
|
||||
@@ -163,6 +165,7 @@ const useStageStoreBase = create<StageState>()((set, get) => ({
|
||||
scenes: [],
|
||||
currentSceneId: null,
|
||||
chats: [],
|
||||
chatSnapshot: { sessions: [], restoreMarker: null },
|
||||
mode: 'playback',
|
||||
toolbarState: 'ai',
|
||||
generatingOutlines: [],
|
||||
@@ -181,6 +184,7 @@ const useStageStoreBase = create<StageState>()((set, get) => ({
|
||||
scenes: [],
|
||||
currentSceneId: null,
|
||||
chats: [],
|
||||
chatSnapshot: { sessions: [], restoreMarker: null },
|
||||
generationComplete: false,
|
||||
generationEpoch: s.generationEpoch + 1,
|
||||
}));
|
||||
@@ -411,7 +415,7 @@ const useStageStoreBase = create<StageState>()((set, get) => ({
|
||||
// durability (e.g. setGenerationComplete) can avoid recording state that
|
||||
// outruns the scene data.
|
||||
saveToStorage: async () => {
|
||||
const { stage, scenes, currentSceneId, chats } = get();
|
||||
const { stage, scenes, currentSceneId, chats, chatSnapshot } = get();
|
||||
if (!stage?.id) {
|
||||
log.warn('Cannot save: stage.id is required');
|
||||
return false;
|
||||
@@ -425,8 +429,21 @@ const useStageStoreBase = create<StageState>()((set, get) => ({
|
||||
scenes: persistedScenes,
|
||||
currentSceneId,
|
||||
chats,
|
||||
chatSnapshot,
|
||||
});
|
||||
|
||||
// Bind future saves to the exact chat snapshot this successful write
|
||||
// represented. Keep the restore marker unchanged: a stale no-op after a
|
||||
// restore must remain stale until the editor reloads.
|
||||
if (get().stage?.id === stage.id && get().chats === chats) {
|
||||
set({
|
||||
chatSnapshot: {
|
||||
sessions: structuredClone(chats),
|
||||
restoreMarker: chatSnapshot.restoreMarker,
|
||||
},
|
||||
});
|
||||
}
|
||||
|
||||
return true;
|
||||
} catch (error) {
|
||||
log.error('Failed to save to storage:', error);
|
||||
@@ -506,6 +523,7 @@ const useStageStoreBase = create<StageState>()((set, get) => ({
|
||||
scenes: migrated,
|
||||
currentSceneId: data.currentSceneId,
|
||||
chats: data.chats,
|
||||
chatSnapshot: data.chatSnapshot ?? { sessions: [], restoreMarker: undefined },
|
||||
outlines,
|
||||
generationComplete,
|
||||
// Compute generatingOutlines from persisted outlines minus completed
|
||||
@@ -540,6 +558,7 @@ const useStageStoreBase = create<StageState>()((set, get) => ({
|
||||
scenes: [],
|
||||
currentSceneId: null,
|
||||
chats: [],
|
||||
chatSnapshot: { sessions: [], restoreMarker: null },
|
||||
outlines: [],
|
||||
generationComplete: false,
|
||||
generationEpoch: s.generationEpoch + 1,
|
||||
|
||||
@@ -54,6 +54,51 @@ export interface ChatSession {
|
||||
lastActionIndex?: number;
|
||||
}
|
||||
|
||||
/**
|
||||
* Advance the session's conflict-order clock without trusting wall time to be
|
||||
* monotonic. Restored data may come from a clock ahead of this device, and the
|
||||
* local clock itself can move backwards.
|
||||
*/
|
||||
export function nextChatUpdatedAt(
|
||||
session: Pick<ChatSession, 'updatedAt'>,
|
||||
now = Date.now(),
|
||||
): number {
|
||||
return Math.max(now, session.updatedAt + 1);
|
||||
}
|
||||
|
||||
/** Apply a lifecycle transition and advance the same conflict-order clock. */
|
||||
export function withChatSessionStatus(
|
||||
session: ChatSession,
|
||||
status: SessionStatus,
|
||||
now = Date.now(),
|
||||
): ChatSession {
|
||||
return { ...session, status, updatedAt: nextChatUpdatedAt(session, now) };
|
||||
}
|
||||
|
||||
/** Advance conflict order once a streamed message segment is fully revealed. */
|
||||
export function withChatSegmentSealed(session: ChatSession, now = Date.now()): ChatSession {
|
||||
return { ...session, updatedAt: nextChatUpdatedAt(session, now) };
|
||||
}
|
||||
|
||||
/** Advance conflict order only when paced text has actually finished revealing. */
|
||||
export function withChatSegmentReveal(
|
||||
session: ChatSession,
|
||||
isComplete: boolean,
|
||||
now = Date.now(),
|
||||
): ChatSession {
|
||||
return isComplete ? withChatSegmentSealed(session, now) : session;
|
||||
}
|
||||
|
||||
/** Mark streams that cannot survive a reload as interrupted without stale ordering. */
|
||||
export function interruptActiveChatSessions(
|
||||
sessions: ChatSession[],
|
||||
now = Date.now(),
|
||||
): ChatSession[] {
|
||||
return sessions.map((session) =>
|
||||
session.status === 'active' ? withChatSessionStatus(session, 'interrupted', now) : session,
|
||||
);
|
||||
}
|
||||
|
||||
/**
|
||||
* Session configuration
|
||||
*/
|
||||
|
||||
@@ -0,0 +1,210 @@
|
||||
const CHAT_STORAGE_GLOBAL_LOCK = 'openmaic:chat-storage:all';
|
||||
const DEFAULT_EXCLUSIVE_ACQUIRE_TIMEOUT_MS = 5_000;
|
||||
type FallbackLockMode = 'shared' | 'exclusive';
|
||||
interface FallbackLockWaiter {
|
||||
mode: FallbackLockMode;
|
||||
start(): void;
|
||||
}
|
||||
|
||||
const fallbackWaiters: FallbackLockWaiter[] = [];
|
||||
let fallbackReaders = 0;
|
||||
let fallbackWriter = false;
|
||||
|
||||
export function chatStoragePartitionLockName(key: string): string {
|
||||
const name = `openmaic:chat-storage:${encodeURIComponent(key)}`;
|
||||
return name === CHAT_STORAGE_GLOBAL_LOCK ? `${name}:partition` : name;
|
||||
}
|
||||
|
||||
function locks(): LockManager | undefined {
|
||||
return typeof navigator !== 'undefined' ? navigator.locks : undefined;
|
||||
}
|
||||
|
||||
function pumpFallbackLocks(): void {
|
||||
if (fallbackWriter || fallbackWaiters.length === 0) return;
|
||||
if (fallbackWaiters[0]!.mode === 'exclusive') {
|
||||
if (fallbackReaders === 0) fallbackWaiters.shift()!.start();
|
||||
return;
|
||||
}
|
||||
while (fallbackWaiters[0]?.mode === 'shared' && !fallbackWriter) {
|
||||
fallbackWaiters.shift()!.start();
|
||||
}
|
||||
}
|
||||
|
||||
function withFallbackRuntimeLock<T>(
|
||||
mode: FallbackLockMode,
|
||||
work: () => Promise<T>,
|
||||
signal?: AbortSignal,
|
||||
): Promise<T> {
|
||||
return new Promise<T>((resolve, reject) => {
|
||||
let started = false;
|
||||
const waiter: FallbackLockWaiter = {
|
||||
mode,
|
||||
start() {
|
||||
started = true;
|
||||
signal?.removeEventListener('abort', onAbort);
|
||||
if (mode === 'shared') fallbackReaders += 1;
|
||||
else fallbackWriter = true;
|
||||
void Promise.resolve()
|
||||
.then(work)
|
||||
.then(resolve, reject)
|
||||
.finally(() => {
|
||||
if (mode === 'shared') fallbackReaders -= 1;
|
||||
else fallbackWriter = false;
|
||||
pumpFallbackLocks();
|
||||
});
|
||||
},
|
||||
};
|
||||
const onAbort = (): void => {
|
||||
if (started) return;
|
||||
const index = fallbackWaiters.indexOf(waiter);
|
||||
if (index >= 0) fallbackWaiters.splice(index, 1);
|
||||
reject(signal?.reason);
|
||||
pumpFallbackLocks();
|
||||
};
|
||||
if (signal?.aborted) {
|
||||
reject(signal.reason);
|
||||
return;
|
||||
}
|
||||
signal?.addEventListener('abort', onAbort, { once: true });
|
||||
fallbackWaiters.push(waiter);
|
||||
pumpFallbackLocks();
|
||||
});
|
||||
}
|
||||
|
||||
/** Let runtime writers run together while excluding whole-store maintenance. */
|
||||
export async function withRuntimeStorageSharedLock<T>(work: () => Promise<T>): Promise<T> {
|
||||
const manager = locks();
|
||||
if (manager) {
|
||||
return manager.request(CHAT_STORAGE_GLOBAL_LOCK, { mode: 'shared' }, work);
|
||||
}
|
||||
return typeof window === 'undefined' ? work() : withFallbackRuntimeLock('shared', work);
|
||||
}
|
||||
|
||||
/**
|
||||
* Bound caller wait time without cancelling the protected work. The shared
|
||||
* lock remains held until `work` really settles, so late storage writes cannot
|
||||
* resume after destructive maintenance has overtaken them.
|
||||
*/
|
||||
export async function withRuntimeStorageSharedLockUntilSettled<T>(
|
||||
work: () => Promise<T>,
|
||||
timeoutMs: number,
|
||||
): Promise<T> {
|
||||
const protectedWork = withRuntimeStorageSharedLock(work);
|
||||
let timer: ReturnType<typeof setTimeout> | undefined;
|
||||
try {
|
||||
return await Promise.race([
|
||||
protectedWork,
|
||||
new Promise<never>((_, reject) => {
|
||||
timer = setTimeout(() => reject(new Error(`timed out after ${timeoutMs}ms`)), timeoutMs);
|
||||
}),
|
||||
]);
|
||||
} finally {
|
||||
clearTimeout(timer);
|
||||
}
|
||||
}
|
||||
|
||||
export interface RuntimeStorageExclusiveLockOptions {
|
||||
acquireTimeoutMs?: number;
|
||||
}
|
||||
|
||||
export class RuntimeStorageLockAcquisitionTimeoutError extends Error {}
|
||||
|
||||
/** Quiesce runtime mutations before destructive whole-store work. */
|
||||
export function withRuntimeStorageExclusiveLock<T>(
|
||||
work: () => Promise<T>,
|
||||
options: RuntimeStorageExclusiveLockOptions = {},
|
||||
): Promise<T> {
|
||||
const manager = locks();
|
||||
if (!manager && typeof window === 'undefined') {
|
||||
return work();
|
||||
}
|
||||
|
||||
const configuredTimeout = options.acquireTimeoutMs ?? DEFAULT_EXCLUSIVE_ACQUIRE_TIMEOUT_MS;
|
||||
const acquireTimeoutMs =
|
||||
Number.isFinite(configuredTimeout) && configuredTimeout > 0
|
||||
? configuredTimeout
|
||||
: DEFAULT_EXCLUSIVE_ACQUIRE_TIMEOUT_MS;
|
||||
let acquired = false;
|
||||
let timer: ReturnType<typeof setTimeout> | undefined;
|
||||
const controller = new AbortController();
|
||||
const timeoutError = new RuntimeStorageLockAcquisitionTimeoutError(
|
||||
`Timed out acquiring the runtime maintenance lock after ${acquireTimeoutMs}ms`,
|
||||
);
|
||||
const guardedWork = async (): Promise<T> => {
|
||||
acquired = true;
|
||||
clearTimeout(timer);
|
||||
return work();
|
||||
};
|
||||
// Cross-realm exclusion is impossible without Web Locks. The fallback still
|
||||
// coordinates every writer in this realm and preserves the pre-cutover
|
||||
// ability to perform an explicit whole-database clear.
|
||||
const request = manager
|
||||
? manager.request(CHAT_STORAGE_GLOBAL_LOCK, { signal: controller.signal }, guardedWork)
|
||||
: withFallbackRuntimeLock('exclusive', guardedWork, controller.signal);
|
||||
|
||||
return new Promise<T>((resolve, reject) => {
|
||||
timer = setTimeout(() => {
|
||||
if (acquired) return;
|
||||
controller.abort(timeoutError);
|
||||
reject(timeoutError);
|
||||
}, acquireTimeoutMs);
|
||||
void request.then(
|
||||
(value) => {
|
||||
clearTimeout(timer);
|
||||
resolve(value);
|
||||
},
|
||||
(error) => {
|
||||
clearTimeout(timer);
|
||||
reject(error);
|
||||
},
|
||||
);
|
||||
});
|
||||
}
|
||||
|
||||
/**
|
||||
* Let a bounded public operation finish while retaining the exclusive lock
|
||||
* until its underlying destructive work actually settles. `releaseCaller`
|
||||
* may be invoked once the caller-visible safety budget has elapsed; the lock
|
||||
* remains owned until `work` itself returns.
|
||||
*/
|
||||
export function withRuntimeStorageExclusiveLockUntilSettled<T>(
|
||||
work: (releaseCaller: (value: T) => void) => Promise<T>,
|
||||
options: RuntimeStorageExclusiveLockOptions = {},
|
||||
): Promise<T> {
|
||||
let callerSettled = false;
|
||||
let resolveCaller!: (value: T) => void;
|
||||
let rejectCaller!: (reason?: unknown) => void;
|
||||
const caller = new Promise<T>((resolve, reject) => {
|
||||
resolveCaller = resolve;
|
||||
rejectCaller = reject;
|
||||
});
|
||||
const releaseCaller = (value: T): void => {
|
||||
if (callerSettled) return;
|
||||
callerSettled = true;
|
||||
resolveCaller(value);
|
||||
};
|
||||
const protectedWork = withRuntimeStorageExclusiveLock(async () => {
|
||||
try {
|
||||
const value = await work(releaseCaller);
|
||||
releaseCaller(value);
|
||||
return value;
|
||||
} catch (error) {
|
||||
if (!callerSettled) {
|
||||
callerSettled = true;
|
||||
rejectCaller(error);
|
||||
}
|
||||
throw error;
|
||||
}
|
||||
}, options);
|
||||
void protectedWork.catch((error) => {
|
||||
if (!callerSettled) {
|
||||
callerSettled = true;
|
||||
rejectCaller(error);
|
||||
}
|
||||
});
|
||||
return caller;
|
||||
}
|
||||
|
||||
/** Compatibility aliases for the chat cutover's partitioned writers. */
|
||||
export const withChatStorageSharedLock = withRuntimeStorageSharedLock;
|
||||
export const withChatStorageExclusiveLock = withRuntimeStorageExclusiveLock;
|
||||
+1616
-49
File diff suppressed because it is too large
Load Diff
+237
-52
@@ -1,4 +1,4 @@
|
||||
import Dexie, { type EntityTable } from 'dexie';
|
||||
import Dexie, { type EntityTable, type Table } from 'dexie';
|
||||
import type {
|
||||
Scene,
|
||||
SceneType,
|
||||
@@ -14,13 +14,21 @@ import type {
|
||||
SessionConfig,
|
||||
ToolCallRecord,
|
||||
ToolCallRequest,
|
||||
ChatSession,
|
||||
} from '@/lib/types/chat';
|
||||
import type { SceneOutline } from '@/lib/types/generation';
|
||||
import type { VoiceDesign } from '@/lib/audio/voice-design';
|
||||
import type { UIMessage } from 'ai';
|
||||
import type { AgentEditSessionRecord } from '@/lib/agent/client/agent-edit-session-types';
|
||||
import { createLogger } from '@/lib/logger';
|
||||
import { deleteStageRuntimeSafely } from '@/lib/runtime/store';
|
||||
import { beginStageRuntimeDeletionSafely, getRuntimeStore } from '@/lib/runtime/store';
|
||||
import type { RuntimeStore } from '@openmaic/storage';
|
||||
import {
|
||||
withRuntimeStorageExclusiveLock,
|
||||
withRuntimeStorageExclusiveLockUntilSettled,
|
||||
withRuntimeStorageSharedLock,
|
||||
} from './chat-storage-lock';
|
||||
import type { ChatStorageOptions } from './chat-storage';
|
||||
|
||||
const log = createLogger('Database');
|
||||
|
||||
@@ -223,7 +231,7 @@ export function mediaFileKey(stageId: string, elementId: string): string {
|
||||
// ==================== Database Definition ====================
|
||||
|
||||
const DATABASE_NAME = 'MAIC-Database';
|
||||
const _DATABASE_VERSION = 12;
|
||||
const _DATABASE_VERSION = 15;
|
||||
|
||||
/**
|
||||
* MAIC Database Instance
|
||||
@@ -236,6 +244,7 @@ class MAICDatabase extends Dexie {
|
||||
imageFiles!: EntityTable<ImageFileRecord, 'id'>;
|
||||
snapshots!: EntityTable<Snapshot, 'id'>; // Undo/redo snapshots (legacy)
|
||||
chatSessions!: EntityTable<ChatSessionRecord, 'id'>;
|
||||
chatRestoreStaging!: Table<ChatSessionRecord, [string, string]>;
|
||||
playbackState!: EntityTable<PlaybackStateRecord, 'stageId'>;
|
||||
stageOutlines!: EntityTable<StageOutlinesRecord, 'stageId'>;
|
||||
mediaFiles!: EntityTable<MediaFileRecord, 'id'>;
|
||||
@@ -443,6 +452,32 @@ class MAICDatabase extends Dexie {
|
||||
autoVoiceCache: 'voiceId, updatedAt',
|
||||
agentEditSessions: 'id, stageId, [stageId+updatedAt]',
|
||||
});
|
||||
|
||||
// Version 13 briefly added chatStorageLocks on the draft chat cutover
|
||||
// branch. Advance past it and remove the abandoned lease table so database
|
||||
// versions stay monotonic for anyone who opened that intermediate build.
|
||||
this.version(14).stores({
|
||||
stages: 'id, updatedAt',
|
||||
scenes: 'id, stageId, order, [stageId+order]',
|
||||
audioFiles: 'id, createdAt',
|
||||
imageFiles: 'id, createdAt',
|
||||
snapshots: '++id',
|
||||
chatSessions: 'id, stageId, [stageId+createdAt]',
|
||||
playbackState: 'stageId',
|
||||
stageOutlines: 'stageId',
|
||||
mediaFiles: 'id, stageId, [stageId+type]',
|
||||
generatedAgents: 'id, stageId',
|
||||
voiceProfiles: 'id, providerId, kind, updatedAt',
|
||||
autoVoiceCache: 'voiceId, updatedAt',
|
||||
agentEditSessions: 'id, stageId, [stageId+updatedAt]',
|
||||
chatStorageLocks: null,
|
||||
});
|
||||
|
||||
// Version 15: backup restore staging must preserve chat IDs that are reused
|
||||
// in different stage partitions; the legacy chat table is keyed by id only.
|
||||
this.version(15).stores({
|
||||
chatRestoreStaging: '[stageId+id], stageId, [stageId+createdAt]',
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
@@ -472,24 +507,80 @@ export async function initDatabase(): Promise<void> {
|
||||
* Clear database (optional)
|
||||
* Use with caution: deletes all data
|
||||
*/
|
||||
export async function clearDatabase(): Promise<void> {
|
||||
await db.delete();
|
||||
export async function clearDatabase(runtimeStore?: RuntimeStore): Promise<void> {
|
||||
// Clear the whole runtime database first, including rows orphaned by an
|
||||
// earlier best-effort stage deletion. This user-requested destructive action
|
||||
// must fail loud: reporting success while runtime data remains is misleading.
|
||||
await withRuntimeStorageExclusiveLock(async () => {
|
||||
await (runtimeStore ?? getRuntimeStore()).deleteAllRuntime();
|
||||
await db.delete();
|
||||
});
|
||||
log.info('Database cleared');
|
||||
}
|
||||
|
||||
function toChatSessionRecord(stageId: string, session: ChatSession): ChatSessionRecord {
|
||||
return {
|
||||
id: session.id,
|
||||
stageId,
|
||||
type: session.type,
|
||||
title: session.title,
|
||||
status: session.status,
|
||||
messages: session.messages,
|
||||
config: session.config,
|
||||
toolCalls: session.toolCalls,
|
||||
pendingToolCalls: session.pendingToolCalls,
|
||||
createdAt: session.createdAt,
|
||||
updatedAt: session.updatedAt,
|
||||
sceneId: session.sceneId,
|
||||
lastActionIndex: session.lastActionIndex,
|
||||
};
|
||||
}
|
||||
|
||||
/**
|
||||
* Export database contents (for backup)
|
||||
*/
|
||||
export async function exportDatabase(): Promise<{
|
||||
export async function exportDatabase(chatOptions: ChatStorageOptions = {}): Promise<{
|
||||
stages: StageRecord[];
|
||||
scenes: SceneRecord[];
|
||||
chatSessions: ChatSessionRecord[];
|
||||
playbackState: PlaybackStateRecord[];
|
||||
}> {
|
||||
const stages = await db.stages.toArray();
|
||||
const legacyChatMap = new Map<string, ChatSessionRecord>();
|
||||
for (const session of [
|
||||
...(await db.chatSessions.toArray()),
|
||||
...(await db.chatRestoreStaging.toArray()),
|
||||
]) {
|
||||
legacyChatMap.set(JSON.stringify([session.stageId, session.id]), session);
|
||||
}
|
||||
const legacyChats = [...legacyChatMap.values()];
|
||||
const { loadChatSessions } = await import('./chat-storage');
|
||||
const runtimeChats = (
|
||||
await Promise.all(
|
||||
stages.map(async (stage) =>
|
||||
(
|
||||
await loadChatSessions(stage.id, {
|
||||
...chatOptions,
|
||||
fallbackToLegacyOnError: false,
|
||||
observe: false,
|
||||
})
|
||||
).map((session) => toChatSessionRecord(stage.id, session)),
|
||||
),
|
||||
)
|
||||
).flat();
|
||||
const runtimeChatKeys = new Set(
|
||||
runtimeChats.map((session) => JSON.stringify([session.stageId, session.id])),
|
||||
);
|
||||
|
||||
return {
|
||||
stages: await db.stages.toArray(),
|
||||
stages,
|
||||
scenes: await db.scenes.toArray(),
|
||||
chatSessions: await db.chatSessions.toArray(),
|
||||
chatSessions: [
|
||||
...runtimeChats,
|
||||
...legacyChats.filter(
|
||||
(session) => !runtimeChatKeys.has(JSON.stringify([session.stageId, session.id])),
|
||||
),
|
||||
],
|
||||
playbackState: await db.playbackState.toArray(),
|
||||
};
|
||||
}
|
||||
@@ -497,23 +588,110 @@ export async function exportDatabase(): Promise<{
|
||||
/**
|
||||
* Import database contents (for restoring backups)
|
||||
*/
|
||||
export async function importDatabase(data: {
|
||||
stages?: StageRecord[];
|
||||
scenes?: SceneRecord[];
|
||||
chatSessions?: ChatSessionRecord[];
|
||||
playbackState?: PlaybackStateRecord[];
|
||||
}): Promise<void> {
|
||||
await db.transaction(
|
||||
'rw',
|
||||
[db.stages, db.scenes, db.chatSessions, db.playbackState],
|
||||
async () => {
|
||||
if (data.stages) await db.stages.bulkPut(data.stages);
|
||||
if (data.scenes) await db.scenes.bulkPut(data.scenes);
|
||||
if (data.chatSessions) await db.chatSessions.bulkPut(data.chatSessions);
|
||||
if (data.playbackState) await db.playbackState.bulkPut(data.playbackState);
|
||||
},
|
||||
);
|
||||
log.info('Database imported successfully');
|
||||
export async function importDatabase(
|
||||
data: {
|
||||
stages?: StageRecord[];
|
||||
scenes?: SceneRecord[];
|
||||
chatSessions?: ChatSessionRecord[];
|
||||
playbackState?: PlaybackStateRecord[];
|
||||
},
|
||||
chatOptions: ChatStorageOptions = {},
|
||||
): Promise<void> {
|
||||
return withRuntimeStorageSharedLock(async () => {
|
||||
const restoredChatStageIds =
|
||||
data.chatSessions === undefined
|
||||
? []
|
||||
: [
|
||||
...new Set([
|
||||
...(data.stages ?? []).map((stage) => stage.id),
|
||||
...data.chatSessions.map((session) => session.stageId),
|
||||
]),
|
||||
];
|
||||
const restoreRows = () =>
|
||||
db.transaction(
|
||||
'rw',
|
||||
[db.stages, db.scenes, db.chatSessions, db.chatRestoreStaging, db.playbackState],
|
||||
async () => {
|
||||
if (data.stages) await db.stages.bulkPut(data.stages);
|
||||
if (data.scenes) await db.scenes.bulkPut(data.scenes);
|
||||
if (data.chatSessions) {
|
||||
for (const stageId of restoredChatStageIds) {
|
||||
await db.chatSessions.where('stageId').equals(stageId).delete();
|
||||
await db.chatRestoreStaging.where('stageId').equals(stageId).delete();
|
||||
}
|
||||
await db.chatRestoreStaging.bulkPut(data.chatSessions);
|
||||
}
|
||||
if (data.playbackState) await db.playbackState.bulkPut(data.playbackState);
|
||||
},
|
||||
);
|
||||
if (data.chatSessions !== undefined) {
|
||||
// RuntimeStore cannot join Dexie's transaction. Keep a rollback image
|
||||
// until durable restore markers have been created for every affected
|
||||
// runtime partition; marker creation failure must not commit half an
|
||||
// imported backup.
|
||||
const rollbackImage = await db.transaction(
|
||||
'r',
|
||||
[db.stages, db.scenes, db.chatSessions, db.chatRestoreStaging, db.playbackState],
|
||||
async () => ({
|
||||
stages: (await db.stages.bulkGet((data.stages ?? []).map((stage) => stage.id))).filter(
|
||||
(stage): stage is StageRecord => stage !== undefined,
|
||||
),
|
||||
scenes: (await db.scenes.bulkGet((data.scenes ?? []).map((scene) => scene.id))).filter(
|
||||
(scene): scene is SceneRecord => scene !== undefined,
|
||||
),
|
||||
chatSessions: (
|
||||
await Promise.all(
|
||||
restoredChatStageIds.map((stageId) =>
|
||||
db.chatSessions.where('stageId').equals(stageId).toArray(),
|
||||
),
|
||||
)
|
||||
).flat(),
|
||||
chatRestoreStaging: (
|
||||
await Promise.all(
|
||||
restoredChatStageIds.map((stageId) =>
|
||||
db.chatRestoreStaging.where('stageId').equals(stageId).toArray(),
|
||||
),
|
||||
)
|
||||
).flat(),
|
||||
playbackState: (
|
||||
await db.playbackState.bulkGet(
|
||||
(data.playbackState ?? []).map((playback) => playback.stageId),
|
||||
)
|
||||
).filter((playback): playback is PlaybackStateRecord => playback !== undefined),
|
||||
}),
|
||||
);
|
||||
const rollbackRows = () =>
|
||||
db.transaction(
|
||||
'rw',
|
||||
[db.stages, db.scenes, db.chatSessions, db.chatRestoreStaging, db.playbackState],
|
||||
async () => {
|
||||
await db.stages.bulkDelete((data.stages ?? []).map((stage) => stage.id));
|
||||
await db.scenes.bulkDelete((data.scenes ?? []).map((scene) => scene.id));
|
||||
for (const stageId of restoredChatStageIds) {
|
||||
await db.chatSessions.where('stageId').equals(stageId).delete();
|
||||
await db.chatRestoreStaging.where('stageId').equals(stageId).delete();
|
||||
}
|
||||
await db.playbackState.bulkDelete(
|
||||
(data.playbackState ?? []).map((playback) => playback.stageId),
|
||||
);
|
||||
await db.stages.bulkPut(rollbackImage.stages);
|
||||
await db.scenes.bulkPut(rollbackImage.scenes);
|
||||
await db.chatSessions.bulkPut(rollbackImage.chatSessions);
|
||||
await db.chatRestoreStaging.bulkPut(rollbackImage.chatRestoreStaging);
|
||||
await db.playbackState.bulkPut(rollbackImage.playbackState);
|
||||
},
|
||||
);
|
||||
const { restoreChatSessionsFromBackup } = await import('./chat-storage');
|
||||
await restoreChatSessionsFromBackup(restoredChatStageIds, restoreRows, {
|
||||
...chatOptions,
|
||||
globalLockHeld: true,
|
||||
rollbackLegacyRows: rollbackRows,
|
||||
});
|
||||
} else {
|
||||
await restoreRows();
|
||||
}
|
||||
log.info('Database imported successfully');
|
||||
});
|
||||
}
|
||||
|
||||
// ==================== Convenience Query Functions ====================
|
||||
@@ -529,33 +707,40 @@ export async function getScenesByStageId(stageId: string): Promise<SceneRecord[]
|
||||
* Delete a course and all its related data
|
||||
*/
|
||||
export async function deleteStageWithRelatedData(stageId: string): Promise<void> {
|
||||
await db.transaction(
|
||||
'rw',
|
||||
[
|
||||
db.stages,
|
||||
db.scenes,
|
||||
db.chatSessions,
|
||||
db.playbackState,
|
||||
db.stageOutlines,
|
||||
db.mediaFiles,
|
||||
db.generatedAgents,
|
||||
db.agentEditSessions,
|
||||
],
|
||||
async () => {
|
||||
await db.stages.delete(stageId);
|
||||
await db.scenes.where('stageId').equals(stageId).delete();
|
||||
await db.chatSessions.where('stageId').equals(stageId).delete();
|
||||
await db.playbackState.delete(stageId);
|
||||
await db.stageOutlines.delete(stageId);
|
||||
await db.mediaFiles.where('stageId').equals(stageId).delete();
|
||||
await db.generatedAgents.where('stageId').equals(stageId).delete();
|
||||
await db.agentEditSessions.where('stageId').equals(stageId).delete();
|
||||
},
|
||||
);
|
||||
// Learner-runtime data lives in a separate IndexedDB database, so it is
|
||||
// cascaded after the Dexie transaction: it cannot join it, and a runtime
|
||||
// failure must not abort it (the helper warns instead of throwing).
|
||||
await deleteStageRuntimeSafely(stageId);
|
||||
await withRuntimeStorageExclusiveLockUntilSettled(async (releaseCaller) => {
|
||||
await db.transaction(
|
||||
'rw',
|
||||
[
|
||||
db.stages,
|
||||
db.scenes,
|
||||
db.chatSessions,
|
||||
db.chatRestoreStaging,
|
||||
db.playbackState,
|
||||
db.stageOutlines,
|
||||
db.mediaFiles,
|
||||
db.generatedAgents,
|
||||
db.agentEditSessions,
|
||||
],
|
||||
async () => {
|
||||
await db.stages.delete(stageId);
|
||||
await db.scenes.where('stageId').equals(stageId).delete();
|
||||
await db.chatSessions.where('stageId').equals(stageId).delete();
|
||||
await db.chatRestoreStaging.where('stageId').equals(stageId).delete();
|
||||
await db.playbackState.delete(stageId);
|
||||
await db.stageOutlines.delete(stageId);
|
||||
await db.mediaFiles.where('stageId').equals(stageId).delete();
|
||||
await db.generatedAgents.where('stageId').equals(stageId).delete();
|
||||
await db.agentEditSessions.where('stageId').equals(stageId).delete();
|
||||
},
|
||||
);
|
||||
// Learner-runtime data lives in a separate IndexedDB database, so it is
|
||||
// cascaded after the Dexie transaction: it cannot join it, and a runtime
|
||||
// failure must not abort it (the helper warns instead of throwing).
|
||||
const runtimeDeletion = beginStageRuntimeDeletionSafely(stageId);
|
||||
await runtimeDeletion.completion;
|
||||
releaseCaller(undefined);
|
||||
await runtimeDeletion.settlement;
|
||||
});
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
+120
-80
@@ -8,12 +8,21 @@
|
||||
import { makeScene, Stage, Scene } from '../types/stage';
|
||||
import { ChatSession } from '../types/chat';
|
||||
import { db } from './database';
|
||||
import { saveChatSessions, loadChatSessions, deleteChatSessions } from './chat-storage';
|
||||
import {
|
||||
saveChatSessions,
|
||||
loadChatSessions,
|
||||
deleteChatSessions,
|
||||
type ChatStorageSnapshot,
|
||||
} from './chat-storage';
|
||||
import { clearPlaybackState } from './playback-storage';
|
||||
import { clearAllForScene } from '@/lib/quiz/persistence';
|
||||
import { deleteStageRuntimeSafely } from '@/lib/runtime/store';
|
||||
import { beginStageRuntimeDeletionSafely } from '@/lib/runtime/store';
|
||||
import { clearStageDrainWatermarks } from '@/lib/pbl/v2/runtime/drain';
|
||||
import { createLogger } from '@/lib/logger';
|
||||
import {
|
||||
withRuntimeStorageExclusiveLockUntilSettled,
|
||||
withRuntimeStorageSharedLock,
|
||||
} from './chat-storage-lock';
|
||||
|
||||
const log = createLogger('StageStorage');
|
||||
|
||||
@@ -22,6 +31,7 @@ export interface StageStoreData {
|
||||
scenes: Scene[];
|
||||
currentSceneId: string | null;
|
||||
chats: ChatSession[];
|
||||
chatSnapshot?: ChatStorageSnapshot;
|
||||
}
|
||||
|
||||
export interface StageListItem {
|
||||
@@ -39,52 +49,57 @@ export interface StageListItem {
|
||||
* Save stage data to IndexedDB
|
||||
*/
|
||||
export async function saveStageData(stageId: string, data: StageStoreData): Promise<void> {
|
||||
try {
|
||||
const now = Date.now();
|
||||
return withRuntimeStorageSharedLock(async () => {
|
||||
try {
|
||||
const now = Date.now();
|
||||
|
||||
// Save to stages table
|
||||
await db.stages.put({
|
||||
id: stageId,
|
||||
name: data.stage.name || 'Untitled Stage',
|
||||
description: data.stage.description,
|
||||
createdAt: data.stage.createdAt || now,
|
||||
updatedAt: now,
|
||||
languageDirective: data.stage.languageDirective,
|
||||
style: data.stage.style,
|
||||
currentSceneId: data.currentSceneId || undefined,
|
||||
agentIds: data.stage.agentIds,
|
||||
videoManifest: data.stage.videoManifest,
|
||||
interactiveMode: data.stage.interactiveMode,
|
||||
taskEngineMode: data.stage.taskEngineMode,
|
||||
generatedAgentConfigs: data.stage.generatedAgentConfigs,
|
||||
});
|
||||
// Save to stages table
|
||||
await db.stages.put({
|
||||
id: stageId,
|
||||
name: data.stage.name || 'Untitled Stage',
|
||||
description: data.stage.description,
|
||||
createdAt: data.stage.createdAt || now,
|
||||
updatedAt: now,
|
||||
languageDirective: data.stage.languageDirective,
|
||||
style: data.stage.style,
|
||||
currentSceneId: data.currentSceneId || undefined,
|
||||
agentIds: data.stage.agentIds,
|
||||
videoManifest: data.stage.videoManifest,
|
||||
interactiveMode: data.stage.interactiveMode,
|
||||
taskEngineMode: data.stage.taskEngineMode,
|
||||
generatedAgentConfigs: data.stage.generatedAgentConfigs,
|
||||
});
|
||||
|
||||
// Delete old scenes first to avoid orphaned data
|
||||
await db.scenes.where('stageId').equals(stageId).delete();
|
||||
// Delete old scenes first to avoid orphaned data
|
||||
await db.scenes.where('stageId').equals(stageId).delete();
|
||||
|
||||
// Save new scenes
|
||||
if (data.scenes && data.scenes.length > 0) {
|
||||
await db.scenes.bulkPut(
|
||||
data.scenes.map((scene, index) => ({
|
||||
...scene,
|
||||
stageId,
|
||||
order: scene.order ?? index,
|
||||
createdAt: scene.createdAt || now,
|
||||
updatedAt: scene.updatedAt || now,
|
||||
})),
|
||||
);
|
||||
// Save new scenes
|
||||
if (data.scenes && data.scenes.length > 0) {
|
||||
await db.scenes.bulkPut(
|
||||
data.scenes.map((scene, index) => ({
|
||||
...scene,
|
||||
stageId,
|
||||
order: scene.order ?? index,
|
||||
createdAt: scene.createdAt || now,
|
||||
updatedAt: scene.updatedAt || now,
|
||||
})),
|
||||
);
|
||||
}
|
||||
|
||||
// Chat sessions live in the learner RuntimeStore, outside the document DB.
|
||||
if (data.chats) {
|
||||
await saveChatSessions(stageId, data.chats, {
|
||||
globalLockHeld: true,
|
||||
snapshot: data.chatSnapshot,
|
||||
});
|
||||
}
|
||||
|
||||
log.info(`Saved stage: ${stageId}`);
|
||||
} catch (error) {
|
||||
log.error('Failed to save stage:', error);
|
||||
throw error;
|
||||
}
|
||||
|
||||
// Save chat sessions to independent table
|
||||
if (data.chats) {
|
||||
await saveChatSessions(stageId, data.chats);
|
||||
}
|
||||
|
||||
log.info(`Saved stage: ${stageId}`);
|
||||
} catch (error) {
|
||||
log.error('Failed to save stage:', error);
|
||||
throw error;
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -102,8 +117,21 @@ export async function loadStageData(stageId: string): Promise<StageStoreData | n
|
||||
// Load scenes
|
||||
const scenes = await db.scenes.where('stageId').equals(stageId).sortBy('order');
|
||||
|
||||
// Load chat sessions from independent table
|
||||
const chats = await loadChatSessions(stageId);
|
||||
// Chat runtime data lives in a separate IndexedDB database. Keep the
|
||||
// document available when that independent store is temporarily
|
||||
// unavailable; a later chat load/save can recover without treating the
|
||||
// already-loaded stage as missing.
|
||||
let chats: ChatSession[] = [];
|
||||
let chatSnapshot: ChatStorageSnapshot = { sessions: [], restoreMarker: undefined };
|
||||
try {
|
||||
chats = await loadChatSessions(stageId, {
|
||||
onSnapshot: (snapshot) => {
|
||||
chatSnapshot = snapshot;
|
||||
},
|
||||
});
|
||||
} catch (error) {
|
||||
log.warn(`Failed to load chat sessions for stage ${stageId}:`, error);
|
||||
}
|
||||
|
||||
log.info(`Loaded stage: ${stageId}, scenes: ${scenes.length}, chats: ${chats.length}`);
|
||||
|
||||
@@ -115,6 +143,7 @@ export async function loadStageData(stageId: string): Promise<StageStoreData | n
|
||||
scenes: scenes.map((s) => makeScene(s, s.content)),
|
||||
currentSceneId: stage.currentSceneId || scenes[0]?.id || null,
|
||||
chats,
|
||||
chatSnapshot,
|
||||
};
|
||||
} catch (error) {
|
||||
log.error('Failed to load stage:', error);
|
||||
@@ -126,42 +155,53 @@ export async function loadStageData(stageId: string): Promise<StageStoreData | n
|
||||
* Delete stage and all related data
|
||||
*/
|
||||
export async function deleteStageData(stageId: string): Promise<void> {
|
||||
try {
|
||||
// Collect scene ids before deletion so we can sweep per-scene localStorage
|
||||
// keys (quiz draft / submitted answers / graded results).
|
||||
const sceneIds = (await db.scenes.where('stageId').equals(stageId).toArray()).map((s) => s.id);
|
||||
|
||||
// Delete stage
|
||||
await db.stages.delete(stageId);
|
||||
|
||||
// Delete scenes
|
||||
await db.scenes.where('stageId').equals(stageId).delete();
|
||||
|
||||
// Delete chat sessions and playback state
|
||||
await deleteChatSessions(stageId);
|
||||
await clearPlaybackState(stageId);
|
||||
|
||||
// Sweep quiz persistence keys for each deleted scene.
|
||||
for (const sceneId of sceneIds) {
|
||||
clearAllForScene(sceneId);
|
||||
}
|
||||
|
||||
// Learner-runtime data lives in a separate IndexedDB database, so it is
|
||||
// cascaded after the Dexie work: it cannot join those transactions, and a
|
||||
// runtime failure must not abort them (the helper warns instead of
|
||||
// throwing).
|
||||
await deleteStageRuntimeSafely(stageId);
|
||||
return withRuntimeStorageExclusiveLockUntilSettled(async (releaseCaller) => {
|
||||
try {
|
||||
await clearStageDrainWatermarks(stageId);
|
||||
} catch (error) {
|
||||
log.warn(`Failed to clear PBL drain watermarks for stage ${stageId}:`, error);
|
||||
}
|
||||
// Collect scene ids before deletion so we can sweep per-scene localStorage
|
||||
// keys (quiz draft / submitted answers / graded results).
|
||||
const sceneIds = (await db.scenes.where('stageId').equals(stageId).toArray()).map(
|
||||
(s) => s.id,
|
||||
);
|
||||
|
||||
log.info(`Deleted stage: ${stageId}`);
|
||||
} catch (error) {
|
||||
log.error('Failed to delete stage:', error);
|
||||
throw error;
|
||||
}
|
||||
// Delete stage
|
||||
await db.stages.delete(stageId);
|
||||
|
||||
// Delete scenes
|
||||
await db.scenes.where('stageId').equals(stageId).delete();
|
||||
|
||||
// Clear legacy chat rows and the still-Dexie playback state. Runtime chat
|
||||
// rows are removed by the all-kind cascade below.
|
||||
await deleteChatSessions(stageId);
|
||||
await clearPlaybackState(stageId);
|
||||
|
||||
// Sweep quiz persistence keys for each deleted scene.
|
||||
for (const sceneId of sceneIds) {
|
||||
clearAllForScene(sceneId);
|
||||
}
|
||||
|
||||
// Learner-runtime data lives in a separate IndexedDB database, so it is
|
||||
// cascaded after the Dexie work: it cannot join those transactions, and a
|
||||
// runtime failure must not abort them (the helper warns instead of
|
||||
// throwing).
|
||||
const runtimeDeletion = beginStageRuntimeDeletionSafely(stageId);
|
||||
await runtimeDeletion.completion;
|
||||
try {
|
||||
await clearStageDrainWatermarks(stageId);
|
||||
} catch (error) {
|
||||
log.warn(`Failed to clear PBL drain watermarks for stage ${stageId}:`, error);
|
||||
}
|
||||
|
||||
log.info(`Deleted stage: ${stageId}`);
|
||||
releaseCaller(undefined);
|
||||
// The public deletion remains bounded, but this callback deliberately
|
||||
// retains the exclusive lock until a late runtime cascade can no longer
|
||||
// delete data written after the caller was released.
|
||||
await runtimeDeletion.settlement;
|
||||
} catch (error) {
|
||||
log.error('Failed to delete stage:', error);
|
||||
throw error;
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -68,6 +68,8 @@ a browser.
|
||||
cascades one learner's sessions + records on one stage, and
|
||||
`deleteStageRuntime` clears a whole stage — the hook a document deletion
|
||||
cascades through.
|
||||
- `deleteAllRuntime` clears every runtime session and record for explicit
|
||||
whole-cache reset flows.
|
||||
|
||||
## Backend equivalence
|
||||
|
||||
|
||||
@@ -444,6 +444,13 @@ export class BrowserRuntimeStore implements RuntimeStore {
|
||||
await this.deleteSessionsByIndex(SESSIONS_BY_STAGE, stageId);
|
||||
}
|
||||
|
||||
async deleteAllRuntime(): Promise<void> {
|
||||
await this.txRun([SESSIONS, RECORDS], 'readwrite', (tx) => {
|
||||
tx.objectStore(SESSIONS).clear();
|
||||
tx.objectStore(RECORDS).clear();
|
||||
});
|
||||
}
|
||||
|
||||
/**
|
||||
* Cascade-delete every session matched by an index query, plus each
|
||||
* session's record range. Idempotent (nothing matched, nothing deleted) and
|
||||
|
||||
@@ -155,4 +155,7 @@ export interface RuntimeStore {
|
||||
* cascades through. Idempotent, deliberately version-agnostic.
|
||||
*/
|
||||
deleteStageRuntime(stageId: string): Promise<void>;
|
||||
|
||||
/** Delete every runtime session and record. Idempotent and version-agnostic. */
|
||||
deleteAllRuntime(): Promise<void>;
|
||||
}
|
||||
|
||||
@@ -315,6 +315,22 @@ export function runRuntimeStoreContract(name: string, makeStore: () => RuntimeSt
|
||||
|
||||
await expect(store.deleteStageRuntime('stage-1')).resolves.toBeUndefined();
|
||||
});
|
||||
|
||||
test('deleteAllRuntime removes every session and record across stages', async () => {
|
||||
const store = makeStore();
|
||||
await store.createSession(makeSession({ id: 's1' }));
|
||||
await store.appendRecord(makeRecordInit('s1'));
|
||||
await store.createSession(makeSession({ id: 's2', stageId: 'stage-2' }));
|
||||
await store.appendRecord(makeRecordInit('s2'));
|
||||
|
||||
await store.deleteAllRuntime();
|
||||
expect(await store.getSession('s1')).toBeUndefined();
|
||||
expect(await store.getSession('s2')).toBeUndefined();
|
||||
expect(await store.listRecords('s1')).toEqual([]);
|
||||
expect(await store.listRecords('s2')).toEqual([]);
|
||||
|
||||
await expect(store.deleteAllRuntime()).resolves.toBeUndefined();
|
||||
});
|
||||
});
|
||||
|
||||
describe('version line', () => {
|
||||
|
||||
@@ -141,6 +141,48 @@ describe('runClassroomLoad', () => {
|
||||
useStageStore.getState().clearStore();
|
||||
});
|
||||
|
||||
it('resets caller-bound chat authority when fallback replaces the classroom', () => {
|
||||
useStageStore.setState({
|
||||
stage: makeStage('stage-old'),
|
||||
chats: [],
|
||||
chatSnapshot: { sessions: [], restoreMarker: 'chat-restore-marker:stage-old:marker' },
|
||||
});
|
||||
|
||||
applyClassroomStageAndScenes(makeStage('stage-new'), [], { persist: false });
|
||||
|
||||
expect(useStageStore.getState().chatSnapshot).toEqual({
|
||||
sessions: [],
|
||||
restoreMarker: null,
|
||||
});
|
||||
useStageStore.getState().clearStore();
|
||||
});
|
||||
|
||||
it('commits runtime chats hydrated for a server fallback', () => {
|
||||
const hydratedChat = {
|
||||
id: 'runtime-chat',
|
||||
type: 'qa' as const,
|
||||
title: 'Runtime chat',
|
||||
status: 'completed' as const,
|
||||
messages: [],
|
||||
config: { agentIds: [] },
|
||||
toolCalls: [],
|
||||
pendingToolCalls: [],
|
||||
createdAt: 1_000,
|
||||
updatedAt: 2_000,
|
||||
};
|
||||
const chatSnapshot = { sessions: [hydratedChat], restoreMarker: null };
|
||||
|
||||
applyClassroomStageAndScenes(makeStage('stage-runtime-chat'), [], {
|
||||
persist: false,
|
||||
chats: [hydratedChat],
|
||||
chatSnapshot,
|
||||
});
|
||||
|
||||
expect(useStageStore.getState().chats).toEqual([hydratedChat]);
|
||||
expect(useStageStore.getState().chatSnapshot).toEqual(chatSnapshot);
|
||||
useStageStore.getState().clearStore();
|
||||
});
|
||||
|
||||
it('does not run stale restore phases after a newer navigation wins', async () => {
|
||||
const loadStorage = deferred<void>();
|
||||
const { deps, setCurrent, setStage } = makeDeps({
|
||||
|
||||
@@ -106,6 +106,7 @@ class MemoryRuntimeStore implements RuntimeStore {
|
||||
|
||||
async deleteLearnerRuntime(): Promise<void> {}
|
||||
async deleteStageRuntime(): Promise<void> {}
|
||||
async deleteAllRuntime(): Promise<void> {}
|
||||
}
|
||||
|
||||
function makeProject(overrides: Partial<PBLProjectV2> = {}): PBLProjectV2 {
|
||||
@@ -194,6 +195,43 @@ afterEach(async () => {
|
||||
});
|
||||
|
||||
describe('classroom server fallback PBL hydration', () => {
|
||||
it('hydrates runtime chats before committing a server fallback', async () => {
|
||||
const stage = makeStage('stage-chat-fallback');
|
||||
const serverScene = makePBLScene(makeProject());
|
||||
const token = claimStageSceneLoadToken();
|
||||
const chatState = {
|
||||
chats: [
|
||||
{
|
||||
id: 'runtime-chat',
|
||||
type: 'qa' as const,
|
||||
title: 'Runtime chat',
|
||||
status: 'completed' as const,
|
||||
messages: [],
|
||||
config: { agentIds: [] },
|
||||
toolCalls: [],
|
||||
pendingToolCalls: [],
|
||||
createdAt: 1_000,
|
||||
updatedAt: 2_000,
|
||||
},
|
||||
],
|
||||
chatSnapshot: { sessions: [], restoreMarker: null },
|
||||
};
|
||||
const applyStageAndScenes = vi.fn();
|
||||
|
||||
await expect(
|
||||
applyHydratedClassroomFallbackScenes({
|
||||
loadToken: token,
|
||||
stage,
|
||||
scenes: [serverScene],
|
||||
hydrateScenes: async () => [serverScene],
|
||||
hydrateChats: async () => chatState,
|
||||
applyStageAndScenes,
|
||||
}),
|
||||
).resolves.toBe(true);
|
||||
|
||||
expect(applyStageAndScenes).toHaveBeenCalledWith(stage, [serverScene], chatState);
|
||||
});
|
||||
|
||||
it('applies fallback scenes under the navigation token that started the request', async () => {
|
||||
const stage = makeStage('stage-a');
|
||||
const serverScene = makePBLScene(makeProject());
|
||||
|
||||
@@ -23,6 +23,7 @@ import {
|
||||
import { advanceMicrotask, startMicrotask } from '@/lib/pbl/v2/operations/progress';
|
||||
import { addSubmission } from '@/lib/pbl/v2/operations/submission';
|
||||
import { clearStageDrainWatermarks, drainProjectRuntime } from '@/lib/pbl/v2/runtime/drain';
|
||||
import { withRuntimeStorageExclusiveLock } from '@/lib/utils/chat-storage-lock';
|
||||
import {
|
||||
enrichPBLRuntimeEvent,
|
||||
PBL_RUNTIME_PAYLOAD_VERSION,
|
||||
@@ -65,6 +66,32 @@ function memoryStorage(): Storage {
|
||||
} as Storage;
|
||||
}
|
||||
|
||||
function serialLockManager(): Pick<LockManager, 'request'> {
|
||||
const tails = new Map<string, Promise<void>>();
|
||||
return {
|
||||
async request<T>(
|
||||
name: string,
|
||||
optionsOrCallback: LockOptions | (() => Promise<T> | T),
|
||||
maybeCallback?: () => Promise<T> | T,
|
||||
): Promise<T> {
|
||||
const callback = typeof optionsOrCallback === 'function' ? optionsOrCallback : maybeCallback!;
|
||||
const previous = tails.get(name) ?? Promise.resolve();
|
||||
let release!: () => void;
|
||||
const current = new Promise<void>((resolve) => {
|
||||
release = resolve;
|
||||
});
|
||||
tails.set(name, current);
|
||||
await previous;
|
||||
try {
|
||||
return await callback();
|
||||
} finally {
|
||||
release();
|
||||
if (tails.get(name) === current) tails.delete(name);
|
||||
}
|
||||
},
|
||||
} as Pick<LockManager, 'request'>;
|
||||
}
|
||||
|
||||
class MemoryKVStore implements KVStore {
|
||||
private readonly values = new Map<string, unknown>();
|
||||
|
||||
@@ -146,6 +173,7 @@ class MemoryRuntimeStore implements RuntimeStore {
|
||||
async deleteLearnerRuntime(): Promise<void> {}
|
||||
|
||||
async deleteStageRuntime(): Promise<void> {}
|
||||
async deleteAllRuntime(): Promise<void> {}
|
||||
}
|
||||
|
||||
class AlreadyExistsRaceStore extends MemoryRuntimeStore {
|
||||
@@ -298,6 +326,7 @@ async function drain(project: PBLProjectV2, store: RuntimeStore, kv: KVStore): P
|
||||
}
|
||||
|
||||
afterEach(() => {
|
||||
vi.unstubAllGlobals();
|
||||
if (vi.isFakeTimers()) {
|
||||
vi.useRealTimers();
|
||||
}
|
||||
@@ -786,6 +815,70 @@ describe('drainProjectRuntime', () => {
|
||||
warn.mockRestore();
|
||||
});
|
||||
|
||||
it('keeps runtime-wide maintenance behind an active PBL drain', async () => {
|
||||
vi.stubGlobal('navigator', { locks: serialLockManager() });
|
||||
const store = new SlowFirstAppendStore('evt-maintenance-lock');
|
||||
const kv = new MemoryKVStore();
|
||||
const draining = drainProjectRuntime({
|
||||
stageId: 'stage-maintenance-lock',
|
||||
sceneId: 'scene-maintenance-lock',
|
||||
project: makeProject([runtimeEvent('evt-maintenance-lock')]),
|
||||
store,
|
||||
kv,
|
||||
learnerKey: LEARNER_KEY,
|
||||
});
|
||||
await vi.waitFor(() => expect(store.appendLog).toEqual(['start:evt-maintenance-lock']));
|
||||
|
||||
let maintenanceStarted = false;
|
||||
const maintenance = withRuntimeStorageExclusiveLock(async () => {
|
||||
maintenanceStarted = true;
|
||||
});
|
||||
await Promise.resolve();
|
||||
expect(maintenanceStarted).toBe(false);
|
||||
|
||||
store.resolveSlowAppend();
|
||||
await draining;
|
||||
await maintenance;
|
||||
expect(maintenanceStarted).toBe(true);
|
||||
});
|
||||
|
||||
it('enrolls queued PBL drains ahead of later runtime maintenance', async () => {
|
||||
vi.stubGlobal('navigator', { locks: serialLockManager() });
|
||||
const store = new SlowFirstAppendStore('evt-first-enrolled');
|
||||
const kv = new MemoryKVStore();
|
||||
const stageId = 'stage-queued-maintenance';
|
||||
const sceneId = 'scene-queued-maintenance';
|
||||
const first = drainProjectRuntime({
|
||||
stageId,
|
||||
sceneId,
|
||||
project: makeProject([runtimeEvent('evt-first-enrolled')]),
|
||||
store,
|
||||
kv,
|
||||
learnerKey: LEARNER_KEY,
|
||||
});
|
||||
await vi.waitFor(() => expect(store.appendLog).toEqual(['start:evt-first-enrolled']));
|
||||
|
||||
const second = drainProjectRuntime({
|
||||
stageId,
|
||||
sceneId,
|
||||
project: makeProject([runtimeEvent('evt-second-enrolled')]),
|
||||
store,
|
||||
kv,
|
||||
learnerKey: LEARNER_KEY,
|
||||
});
|
||||
await Promise.resolve();
|
||||
const maintenance = withRuntimeStorageExclusiveLock(async () => {
|
||||
store.records.splice(0);
|
||||
store.sessions.splice(0);
|
||||
});
|
||||
|
||||
store.resolveSlowAppend();
|
||||
await Promise.all([first, second, maintenance]);
|
||||
|
||||
expect(store.records).toEqual([]);
|
||||
expect(store.sessions).toEqual([]);
|
||||
});
|
||||
|
||||
it('does not permanently block later drains after a queued append times out', async () => {
|
||||
vi.useFakeTimers();
|
||||
const warn = vi.spyOn(console, 'warn').mockImplementation(() => {});
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
import { describe, expect, it, vi } from 'vitest';
|
||||
import { afterEach, describe, expect, it, vi } from 'vitest';
|
||||
import { IDBFactory, IDBKeyRange } from 'fake-indexeddb';
|
||||
import type {
|
||||
RuntimePayload,
|
||||
@@ -40,6 +40,7 @@ import {
|
||||
import { makeScene, type Scene } from '@/lib/types/stage';
|
||||
import type { KVScope, KVStore } from '@openmaic/storage';
|
||||
import type { PBLProjectV2, PBLRuntimeEvent } from '@/lib/pbl/v2/types';
|
||||
import { withRuntimeStorageExclusiveLock } from '@/lib/utils/chat-storage-lock';
|
||||
|
||||
if (!('IDBKeyRange' in globalThis)) {
|
||||
Object.defineProperty(globalThis, 'IDBKeyRange', { value: IDBKeyRange, configurable: true });
|
||||
@@ -49,6 +50,68 @@ const STAGE_ID = 'stage-1';
|
||||
const SCENE_ID = 'scene-1';
|
||||
const LEARNER_KEY = 'anon:test-device';
|
||||
|
||||
function readWriteLockManager(onAcquire?: (mode: LockMode) => void): Pick<LockManager, 'request'> {
|
||||
type Waiter = {
|
||||
mode: LockMode;
|
||||
run(): void;
|
||||
};
|
||||
const states = new Map<string, { readers: number; writer: boolean; waiters: Waiter[] }>();
|
||||
const stateFor = (name: string) => {
|
||||
let state = states.get(name);
|
||||
if (!state) {
|
||||
state = { readers: 0, writer: false, waiters: [] };
|
||||
states.set(name, state);
|
||||
}
|
||||
return state;
|
||||
};
|
||||
const pump = (name: string) => {
|
||||
const state = stateFor(name);
|
||||
if (state.writer || state.waiters.length === 0) return;
|
||||
if (state.waiters[0]!.mode === 'exclusive') {
|
||||
if (state.readers === 0) state.waiters.shift()!.run();
|
||||
return;
|
||||
}
|
||||
while (state.waiters[0]?.mode === 'shared' && !state.writer) {
|
||||
state.waiters.shift()!.run();
|
||||
}
|
||||
};
|
||||
return {
|
||||
request<T>(
|
||||
name: string,
|
||||
optionsOrCallback: LockOptions | (() => Promise<T> | T),
|
||||
maybeCallback?: () => Promise<T> | T,
|
||||
): Promise<T> {
|
||||
const options = typeof optionsOrCallback === 'function' ? undefined : optionsOrCallback;
|
||||
const callback = typeof optionsOrCallback === 'function' ? optionsOrCallback : maybeCallback!;
|
||||
const mode = options?.mode ?? 'exclusive';
|
||||
const state = stateFor(name);
|
||||
return new Promise<T>((resolve, reject) => {
|
||||
state.waiters.push({
|
||||
mode,
|
||||
run() {
|
||||
if (mode === 'shared') state.readers += 1;
|
||||
else state.writer = true;
|
||||
onAcquire?.(mode);
|
||||
void Promise.resolve()
|
||||
.then(callback)
|
||||
.then(resolve, reject)
|
||||
.finally(() => {
|
||||
if (mode === 'shared') state.readers -= 1;
|
||||
else state.writer = false;
|
||||
pump(name);
|
||||
});
|
||||
},
|
||||
});
|
||||
pump(name);
|
||||
});
|
||||
},
|
||||
} as Pick<LockManager, 'request'>;
|
||||
}
|
||||
|
||||
afterEach(() => {
|
||||
vi.unstubAllGlobals();
|
||||
});
|
||||
|
||||
class MemoryKVStore implements KVStore {
|
||||
private readonly values = new Map<string, unknown>();
|
||||
|
||||
@@ -120,6 +183,7 @@ class MemoryRuntimeStore implements RuntimeStore {
|
||||
|
||||
async deleteLearnerRuntime(): Promise<void> {}
|
||||
async deleteStageRuntime(): Promise<void> {}
|
||||
async deleteAllRuntime(): Promise<void> {}
|
||||
}
|
||||
|
||||
class ThrowingRuntimeStore extends MemoryRuntimeStore {
|
||||
@@ -614,6 +678,49 @@ describe('PBL runtime hydration', () => {
|
||||
}
|
||||
});
|
||||
|
||||
it('enrolls hydration before later maintenance while waiting for an earlier drain', async () => {
|
||||
let sharedAcquisitions = 0;
|
||||
vi.stubGlobal('navigator', {
|
||||
locks: readWriteLockManager((mode) => {
|
||||
if (mode === 'shared') sharedAcquisitions += 1;
|
||||
}),
|
||||
});
|
||||
const project = transitionProjectUiPhase(makeProject(), 'workspace');
|
||||
const eventId = project.runtimeEvents![0]!.id;
|
||||
const store = new SlowFirstAppendRuntimeStore(eventId);
|
||||
const kv = new MemoryKVStore();
|
||||
|
||||
const priorDrain = drainProjectRuntime({
|
||||
stageId: STAGE_ID,
|
||||
sceneId: SCENE_ID,
|
||||
project,
|
||||
store,
|
||||
kv,
|
||||
learnerKey: LEARNER_KEY,
|
||||
});
|
||||
await store.appendStarted;
|
||||
const hydrating = hydratePBLProjectFromRuntime({
|
||||
stageId: STAGE_ID,
|
||||
sceneId: SCENE_ID,
|
||||
project,
|
||||
store,
|
||||
kv,
|
||||
learnerKey: LEARNER_KEY,
|
||||
});
|
||||
await new Promise((resolve) => setTimeout(resolve, 0));
|
||||
const sharedAcquisitionsBeforeRelease = sharedAcquisitions;
|
||||
const maintenance = withRuntimeStorageExclusiveLock(async () => {
|
||||
store.records.splice(0);
|
||||
store.sessions.splice(0);
|
||||
});
|
||||
|
||||
store.release();
|
||||
await Promise.all([priorDrain, hydrating, maintenance]);
|
||||
expect(sharedAcquisitionsBeforeRelease).toBe(2);
|
||||
expect(store.records).toEqual([]);
|
||||
expect(store.sessions).toEqual([]);
|
||||
});
|
||||
|
||||
it('falls back to the document without a snapshot when the drain barrier times out', async () => {
|
||||
vi.useFakeTimers();
|
||||
const warn = vi.spyOn(console, 'warn').mockImplementation(() => {});
|
||||
@@ -665,6 +772,45 @@ describe('PBL runtime hydration', () => {
|
||||
}
|
||||
});
|
||||
|
||||
it('retains the shared maintenance lock after a timed-out hydration until work settles', async () => {
|
||||
vi.useFakeTimers();
|
||||
vi.stubGlobal('navigator', { locks: readWriteLockManager() });
|
||||
const project = transitionProjectUiPhase(makeProject(), 'workspace');
|
||||
const eventId = project.runtimeEvents![0]!.id;
|
||||
const store = new SlowFirstAppendRuntimeStore(eventId);
|
||||
|
||||
try {
|
||||
const hydrating = hydrate(project, store);
|
||||
let hydrationSettled = false;
|
||||
const hydrationOutcome = hydrating.then(
|
||||
() => {
|
||||
hydrationSettled = true;
|
||||
},
|
||||
() => {
|
||||
hydrationSettled = true;
|
||||
},
|
||||
);
|
||||
await store.appendStarted;
|
||||
await vi.advanceTimersByTimeAsync(20_001);
|
||||
await Promise.resolve();
|
||||
expect(hydrationSettled).toBe(true);
|
||||
|
||||
let maintenanceStarted = false;
|
||||
const maintenance = withRuntimeStorageExclusiveLock(async () => {
|
||||
maintenanceStarted = true;
|
||||
});
|
||||
await Promise.resolve();
|
||||
expect(maintenanceStarted).toBe(false);
|
||||
|
||||
store.release();
|
||||
await Promise.all([hydrationOutcome, maintenance]);
|
||||
expect(maintenanceStarted).toBe(true);
|
||||
} finally {
|
||||
store.release();
|
||||
vi.useRealTimers();
|
||||
}
|
||||
});
|
||||
|
||||
it('keeps document state and self-heals when runtime history is partial', async () => {
|
||||
const warn = vi.spyOn(console, 'warn').mockImplementation(() => {});
|
||||
const store = new MemoryRuntimeStore();
|
||||
|
||||
File diff suppressed because it is too large
Load Diff
File diff suppressed because it is too large
Load Diff
@@ -5,7 +5,10 @@ import { describe, expect, it, vi } from 'vitest';
|
||||
// established pattern for database-touching code — no Dexie-in-node harness)
|
||||
// and run the REAL function to prove it cascades into the runtime layer.
|
||||
vi.mock('@/lib/runtime/store', () => ({
|
||||
deleteStageRuntimeSafely: vi.fn().mockResolvedValue(undefined),
|
||||
beginStageRuntimeDeletionSafely: vi.fn(() => ({
|
||||
completion: Promise.resolve(),
|
||||
settlement: Promise.resolve(),
|
||||
})),
|
||||
}));
|
||||
vi.mock('@/lib/utils/database', () => ({
|
||||
db: {
|
||||
@@ -31,13 +34,20 @@ vi.mock('@/lib/utils/playback-storage', () => ({
|
||||
vi.mock('@/lib/quiz/persistence', () => ({
|
||||
clearAllForScene: vi.fn(),
|
||||
}));
|
||||
vi.mock('@/lib/utils/chat-storage-lock', () => ({
|
||||
withRuntimeStorageExclusiveLockUntilSettled: vi.fn(
|
||||
async (work: (releaseCaller: (value: unknown) => void) => Promise<unknown>) => work(() => {}),
|
||||
),
|
||||
}));
|
||||
|
||||
import { deleteStageData } from '@/lib/utils/stage-storage';
|
||||
import { deleteStageRuntimeSafely } from '@/lib/runtime/store';
|
||||
import { beginStageRuntimeDeletionSafely } from '@/lib/runtime/store';
|
||||
import { withRuntimeStorageExclusiveLockUntilSettled } from '@/lib/utils/chat-storage-lock';
|
||||
|
||||
describe('deleteStageData runtime cascade', () => {
|
||||
it('cascades into the runtime store with the deleted stageId', async () => {
|
||||
await deleteStageData('stage-7');
|
||||
expect(vi.mocked(deleteStageRuntimeSafely)).toHaveBeenCalledExactlyOnceWith('stage-7');
|
||||
expect(vi.mocked(beginStageRuntimeDeletionSafely)).toHaveBeenCalledExactlyOnceWith('stage-7');
|
||||
expect(vi.mocked(withRuntimeStorageExclusiveLockUntilSettled)).toHaveBeenCalledOnce();
|
||||
});
|
||||
});
|
||||
|
||||
@@ -0,0 +1,68 @@
|
||||
import { describe, expect, it, vi } from 'vitest';
|
||||
|
||||
const { loadChatSessions } = vi.hoisted(() => ({
|
||||
loadChatSessions: vi.fn().mockRejectedValue(new Error('runtime unavailable')),
|
||||
}));
|
||||
|
||||
vi.mock('@/lib/utils/database', () => ({
|
||||
db: {
|
||||
stages: {
|
||||
get: vi.fn().mockResolvedValue({
|
||||
id: 'stage-1',
|
||||
name: 'Persisted stage',
|
||||
createdAt: 1_000,
|
||||
updatedAt: 2_000,
|
||||
}),
|
||||
},
|
||||
scenes: {
|
||||
where: () => ({
|
||||
equals: () => ({
|
||||
sortBy: vi.fn().mockResolvedValue([
|
||||
{
|
||||
id: 'scene-1',
|
||||
stageId: 'stage-1',
|
||||
title: 'Persisted scene',
|
||||
order: 0,
|
||||
content: { type: 'slide', canvas: {} },
|
||||
},
|
||||
]),
|
||||
}),
|
||||
}),
|
||||
},
|
||||
},
|
||||
}));
|
||||
vi.mock('@/lib/utils/chat-storage', () => ({
|
||||
saveChatSessions: vi.fn(),
|
||||
loadChatSessions,
|
||||
deleteChatSessions: vi.fn(),
|
||||
}));
|
||||
vi.mock('@/lib/utils/playback-storage', () => ({
|
||||
clearPlaybackState: vi.fn(),
|
||||
}));
|
||||
vi.mock('@/lib/quiz/persistence', () => ({
|
||||
clearAllForScene: vi.fn(),
|
||||
}));
|
||||
vi.mock('@/lib/runtime/store', () => ({
|
||||
beginStageRuntimeDeletionSafely: vi.fn(),
|
||||
}));
|
||||
vi.mock('@/lib/pbl/v2/runtime/drain', () => ({
|
||||
clearStageDrainWatermarks: vi.fn(),
|
||||
}));
|
||||
|
||||
import { loadStageData } from '@/lib/utils/stage-storage';
|
||||
|
||||
describe('loadStageData chat failure isolation', () => {
|
||||
it('keeps persisted stage and scene data available when chat storage fails', async () => {
|
||||
await expect(loadStageData('stage-1')).resolves.toMatchObject({
|
||||
stage: { id: 'stage-1', name: 'Persisted stage' },
|
||||
scenes: [{ id: 'scene-1', type: 'slide', title: 'Persisted scene' }],
|
||||
currentSceneId: 'scene-1',
|
||||
chats: [],
|
||||
chatSnapshot: { sessions: [], restoreMarker: undefined },
|
||||
});
|
||||
expect(loadChatSessions).toHaveBeenCalledExactlyOnceWith(
|
||||
'stage-1',
|
||||
expect.objectContaining({ onSnapshot: expect.any(Function) }),
|
||||
);
|
||||
});
|
||||
});
|
||||
Reference in New Issue
Block a user