From 05ab9f9f077331ab81aa08cf39a5235bb2f90b78 Mon Sep 17 00:00:00 2001 From: ttbombadil Date: Sun, 4 Oct 2026 15:47:02 +0200 Subject: [PATCH 1/5] fix(tutor): propagate request cancellation --- client/src/hooks/use-tutor.ts | 29 +++++- server/routes/tutor.routes.ts | 43 ++++++++- .../course-content/course-content-loader.ts | 5 +- server/services/examples/http-provider.ts | 5 +- server/services/examples/source-provider.ts | 1 + server/services/tutor/kiconnect-provider.ts | 87 ++++++++++-------- server/services/tutor/llm-provider.ts | 4 +- server/services/tutor/tutor-service.ts | 88 +++++++++++++++---- tests/client/tutor.test.tsx | 32 +++++++ tests/server/examples/http-provider.test.ts | 14 +++ tests/server/examples/source-provider.test.ts | 24 ++++- tests/server/routes/tutor.routes.test.ts | 61 +++++++++++-- .../course-content-loader.test.ts | 20 +++++ .../services/tutor/kiconnect-provider.test.ts | 48 ++++++++++ .../tutor/session-concurrency.test.ts | 57 ++++++++++++ 15 files changed, 446 insertions(+), 72 deletions(-) diff --git a/client/src/hooks/use-tutor.ts b/client/src/hooks/use-tutor.ts index caa45dc9f..7eecb9d66 100644 --- a/client/src/hooks/use-tutor.ts +++ b/client/src/hooks/use-tutor.ts @@ -111,6 +111,10 @@ function getCourseContentRequest( return {}; } +function isAbortError(error: unknown): boolean { + return error instanceof Error && error.name === "AbortError"; +} + function buildDialogTurn( question: string, answer: string, @@ -227,15 +231,23 @@ export function useTutor( const [isLoading, setIsLoading] = useState(false); const [error, setError] = useState(null); const canLoadConfigOnMount = useRef(capabilities.canUseTutor); + const activeTutorRequests = useRef(new Set()); + const abortPendingTutorRequests = useCallback(() => { + for (const controller of activeTutorRequests.current) controller.abort(); + activeTutorRequests.current.clear(); + }, []); useEffect(() => { + abortPendingTutorRequests(); setCourseContentSession(undefined); setHistory([]); setQuestion(null); setAnswer(""); setError(null); setEffectiveDifficulty(configuredDifficulty); - }, [courseContentKey]); + }, [abortPendingTutorRequests, courseContentKey]); + + useEffect(() => () => abortPendingTutorRequests(), [abortPendingTutorRequests]); useEffect(() => { if (!canLoadConfigOnMount.current) return; @@ -274,12 +286,15 @@ export function useTutor( } setModelsLoading(true); + const controller = new AbortController(); + activeTutorRequests.current.add(controller); try { const response = await fetch("/api/tutor/models", { method: "POST", credentials: "include", headers: { "Content-Type": "application/json" }, body: JSON.stringify({ credential }), + signal: controller.signal, }); const body: unknown = await response.json().catch(() => null); if (!response.ok) throw new Error(getErrorMessage(body)); @@ -288,10 +303,12 @@ export function useTutor( setAvailableModels(parsed.data.models); setSelectedModel((current) => current === "auto" || parsed.data.models.includes(current) ? current : "auto"); } catch (requestError) { + if (controller.signal.aborted || isAbortError(requestError)) return; setAvailableModels([]); setSelectedModel("auto"); setError(requestError instanceof Error ? requestError.message : "The Tutor models could not be loaded."); } finally { + activeTutorRequests.current.delete(controller); setModelsLoading(false); } }, [capabilities, credential]); @@ -311,11 +328,14 @@ export function useTutor( if (selectedModel !== "auto" && !requestedModel) setSelectedModel("auto"); setIsLoading(true); + const controller = new AbortController(); + activeTutorRequests.current.add(controller); try { const response = await fetch("/api/tutor/question", { method: "POST", credentials: "include", headers: { "Content-Type": "application/json" }, + signal: controller.signal, body: JSON.stringify({ code, credential, @@ -335,8 +355,10 @@ export function useTutor( setLastUsedModel(parsed.data.model); setAnswer(""); } catch (requestError) { + if (controller.signal.aborted || isAbortError(requestError)) return; setError(requestError instanceof Error ? requestError.message : "The Tutor request failed."); } finally { + activeTutorRequests.current.delete(controller); setIsLoading(false); } }, [availableModels, capabilities, courseContent, courseContentSession, credential, selectedModel]); @@ -381,11 +403,14 @@ export function useTutor( const currentHistory = history; setIsLoading(true); + const controller = new AbortController(); + activeTutorRequests.current.add(controller); try { const response = await fetch("/api/tutor/dialog", { method: "POST", credentials: "include", headers: { "Content-Type": "application/json" }, + signal: controller.signal, body: JSON.stringify({ code, history: currentHistory, @@ -413,8 +438,10 @@ export function useTutor( } catch (requestError) { // Keep the current question, answer, and history intact so a failed // request can be retried deliberately by the learner. + if (controller.signal.aborted || isAbortError(requestError)) return; setError(requestError instanceof Error ? requestError.message : "The Tutor request failed."); } finally { + activeTutorRequests.current.delete(controller); setIsLoading(false); } }, [answer, availableModels, capabilities, courseContent, courseContentSession, credential, effectiveDifficulty, history, question, selectedModel]); diff --git a/server/routes/tutor.routes.ts b/server/routes/tutor.routes.ts index 23b7469cf..14acab594 100644 --- a/server/routes/tutor.routes.ts +++ b/server/routes/tutor.routes.ts @@ -103,6 +103,25 @@ function requireRequestCredential(req: Request, res: Response, credential: strin return true; } +function bindRequestAbortSignal(req: Request, res: Response): AbortSignal { + const controller = new AbortController(); + const abort = () => controller.abort(new DOMException("Request aborted", "AbortError")); + const cleanup = () => { + req.off("aborted", abort); + res.off("close", onClose); + res.off("finish", cleanup); + }; + const onClose = () => { + if (!res.writableEnded) abort(); + cleanup(); + }; + req.once("aborted", abort); + res.once("close", onClose); + res.once("finish", cleanup); + if (req.aborted || res.destroyed) abort(); + return controller.signal; +} + export function registerTutorRoutes(app: Express, deps: TutorRouteDeps = {}): void { const logger = deps.logger ?? new Logger("TutorRoutes"); const service = deps.service ?? createTutorService(new KiconnectProvider()); @@ -110,6 +129,7 @@ export function registerTutorRoutes(app: Express, deps: TutorRouteDeps = {}): vo const sessionStore = deps.sessionStore ?? new TutorCourseContentSessionStore(); app.post("/api/tutor/question", async (req, res) => { + const signal = bindRequestAbortSignal(req, res); if (!enforceTutorRateLimit(res, { ...deps, rateLimiter })) return; const parsed = tutorQuestionRequestSchema.safeParse(req.body); @@ -120,7 +140,8 @@ export function registerTutorRoutes(app: Express, deps: TutorRouteDeps = {}): vo if (!requireRequestCredential(req, res, parsed.data.credential)) return; try { - const content = await resolveTutorContent(req, res, parsed.data.courseContent, parsed.data.courseContentSession, deps.courseContent, sessionStore); + const content = await resolveTutorContent(req, res, parsed.data.courseContent, parsed.data.courseContentSession, deps.courseContent, sessionStore, signal); + if (signal.aborted) return; if ((parsed.data.courseContent !== undefined || parsed.data.courseContentSession !== undefined) && !content) return; const generated = await service.generateQuestion( parsed.data.code, @@ -128,7 +149,9 @@ export function registerTutorRoutes(app: Express, deps: TutorRouteDeps = {}): vo parsed.data.model, parsed.data.difficulty, content?.planning, + signal, ); + if (signal.aborted) return; res.json({ ...generated.result, provider: config.tutor.provider, @@ -136,6 +159,7 @@ export function registerTutorRoutes(app: Express, deps: TutorRouteDeps = {}): vo ...(content?.session ? { courseContentSession: content.session } : {}), }); } catch (error) { + if (signal.aborted) return; if (error instanceof TutorProviderError) { mapProviderError(res, error); return; @@ -146,6 +170,7 @@ export function registerTutorRoutes(app: Express, deps: TutorRouteDeps = {}): vo }); app.post("/api/tutor/models", async (req, res) => { + const signal = bindRequestAbortSignal(req, res); if (!enforceTutorRateLimit(res, { ...deps, rateLimiter })) return; const parsed = tutorModelsRequestSchema.safeParse(req.body); @@ -156,9 +181,11 @@ export function registerTutorRoutes(app: Express, deps: TutorRouteDeps = {}): vo if (!requireRequestCredential(req, res, parsed.data.credential)) return; try { - const models = await service.getAvailableModels(parsed.data.credential); + const models = await service.getAvailableModels(parsed.data.credential, signal); + if (signal.aborted) return; res.json({ models }); } catch (error) { + if (signal.aborted) return; if (error instanceof TutorProviderError) { mapProviderError(res, error); return; @@ -169,6 +196,7 @@ export function registerTutorRoutes(app: Express, deps: TutorRouteDeps = {}): vo }); app.post("/api/tutor/dialog", async (req, res) => { + const signal = bindRequestAbortSignal(req, res); if (!enforceTutorRateLimit(res, { ...deps, rateLimiter })) return; const parsed = tutorDialogRequestSchema.safeParse(req.body); @@ -179,7 +207,8 @@ export function registerTutorRoutes(app: Express, deps: TutorRouteDeps = {}): vo if (!requireRequestCredential(req, res, parsed.data.credential)) return; try { - const content = await resolveTutorContent(req, res, parsed.data.courseContent, parsed.data.courseContentSession, deps.courseContent, sessionStore); + const content = await resolveTutorContent(req, res, parsed.data.courseContent, parsed.data.courseContentSession, deps.courseContent, sessionStore, signal); + if (signal.aborted) return; if ((parsed.data.courseContent !== undefined || parsed.data.courseContentSession !== undefined) && !content) return; const generated = await service.generateDialogResponse( parsed.data.code, @@ -190,7 +219,9 @@ export function registerTutorRoutes(app: Express, deps: TutorRouteDeps = {}): vo parsed.data.model, parsed.data.difficulty, content?.planning, + signal, ); + if (signal.aborted) return; res.json({ ...generated.result, provider: config.tutor.provider, @@ -198,6 +229,7 @@ export function registerTutorRoutes(app: Express, deps: TutorRouteDeps = {}): vo ...(content?.session ? { courseContentSession: content.session } : {}), }); } catch (error) { + if (signal.aborted) return; if (error instanceof TutorProviderError) { mapProviderError(res, error); return; @@ -215,7 +247,9 @@ async function resolveTutorContent( sessionHandle: string | undefined, resolver: TutorCourseContentResolver | undefined, sessions: TutorCourseContentSessionStore, + signal: AbortSignal, ): Promise<{ readonly planning: TutorPlanningContentContext; readonly session: string } | undefined> { + if (signal.aborted) return undefined; const identity = (res.locals.unosimIdentity as RequestIdentity | undefined)?.subject ?? "anonymous"; if (sessionHandle !== undefined) { const pinned = sessions.get(identity, sessionHandle); @@ -233,14 +267,17 @@ async function resolveTutorContent( const context: RequestContext = { identity, requestId: req.header("x-request-id") ?? randomUUID(), + signal, }; let resolved: ResolvedTutorCourseContent; try { resolved = await resolver.resolveTutorContent(request, context); } catch { + if (signal.aborted) return undefined; responseError(res, 400, "INVALID_REQUEST", "Der Course-Content-Kontext ist ungültig oder nicht mehr aktiv."); return undefined; } + if (signal.aborted) return undefined; const handle = sessions.create(identity, resolved); return { planning: resolved, session: handle }; } diff --git a/server/services/course-content/course-content-loader.ts b/server/services/course-content/course-content-loader.ts index 32c7fc4c6..5016d57ee 100644 --- a/server/services/course-content/course-content-loader.ts +++ b/server/services/course-content/course-content-loader.ts @@ -141,7 +141,10 @@ export class CourseContentLoader { const strategies = await this.loadStrategies(base, manifest, signal); validateTutorReferences(manifest, topicEntries, strategies, annotations); return { status: "valid", manifest, topics: topicEntries, strategies }; - } catch { + } catch (error) { + if (signal?.aborted || (error instanceof Error && error.name === "AbortError")) { + throw signal?.reason ?? error; + } return { status: "invalid", reason: "Tutor capability is invalid" }; } } diff --git a/server/services/examples/http-provider.ts b/server/services/examples/http-provider.ts index c5e7e6961..97ebbfaaa 100644 --- a/server/services/examples/http-provider.ts +++ b/server/services/examples/http-provider.ts @@ -51,11 +51,14 @@ export function validateSourceUrl(value: string): URL { } async function fetchText(url: URL, maxBytes: number, requestSignal?: AbortSignal): Promise { + requestSignal?.throwIfAborted(); const validated = validateSourceUrl(url.toString()); await assertPublicHost(validated.hostname); + requestSignal?.throwIfAborted(); const controller = new AbortController(); const abort = () => controller.abort(requestSignal?.reason); - requestSignal?.addEventListener("abort", abort, { once: true }); + if (requestSignal?.aborted) abort(); + else requestSignal?.addEventListener("abort", abort, { once: true }); const timeout = setTimeout(() => controller.abort(), config.examples.timeoutMs); try { const response = await fetch(validated, { signal: controller.signal, redirect: "manual" }); diff --git a/server/services/examples/source-provider.ts b/server/services/examples/source-provider.ts index 253f7e705..722930eb7 100644 --- a/server/services/examples/source-provider.ts +++ b/server/services/examples/source-provider.ts @@ -94,6 +94,7 @@ export class SourceProvider { }, context.signal), ); } catch (error) { + if (context.signal?.aborted) throw context.signal.reason ?? error; const lastGood = this.cache.markSourceFailure( sourceKey, this.now() + Math.min( diff --git a/server/services/tutor/kiconnect-provider.ts b/server/services/tutor/kiconnect-provider.ts index 316662dca..6d01f169c 100644 --- a/server/services/tutor/kiconnect-provider.ts +++ b/server/services/tutor/kiconnect-provider.ts @@ -15,6 +15,39 @@ import { rankTutorModels } from "./model-preference"; export const TUTOR_TEMPERATURE = 0.2; +function isAbortError(error: unknown): boolean { + return error instanceof Error && error.name === "AbortError"; +} + +async function withProviderTimeout( + requestSignal: AbortSignal | undefined, + operation: (signal: AbortSignal) => Promise, +): Promise { + requestSignal?.throwIfAborted(); + const controller = new AbortController(); + let timedOut = false; + const abortFromRequest = () => controller.abort(requestSignal?.reason); + requestSignal?.addEventListener("abort", abortFromRequest, { once: true }); + const timeout = globalThis.setTimeout(() => { + timedOut = true; + controller.abort(); + }, config.tutor.timeoutMs); + + try { + const result = await operation(controller.signal); + requestSignal?.throwIfAborted(); + return result; + } catch (error) { + if (requestSignal?.aborted) throw requestSignal.reason ?? error; + if (error instanceof TutorProviderError) throw error; + if (timedOut || isAbortError(error)) throw new TutorProviderError("provider-timeout"); + throw new TutorProviderError("provider-unavailable"); + } finally { + globalThis.clearTimeout(timeout); + requestSignal?.removeEventListener("abort", abortFromRequest); + } +} + const completionChoicesSchema = z.array(z.object({ message: z.object({ content: z.unknown(), @@ -190,13 +223,11 @@ export class KiconnectProvider implements LLMProvider, StructuredLLMProvider { private readonly baseUrl = config.tutor.baseUrl; private readonly logger = new Logger("KiconnectProvider"); - async listModels(credential: string): Promise { - const controller = new AbortController(); - const timeout = globalThis.setTimeout(() => controller.abort(), config.tutor.timeoutMs); - try { + async listModels(credential: string, signal?: AbortSignal): Promise { + return withProviderTimeout(signal, async (providerSignal) => { const response = await fetch(`${this.baseUrl}/models`, { headers: { Authorization: `Bearer ${credential}` }, - signal: controller.signal, + signal: providerSignal, }); if (!response.ok) { throw providerErrorForStatus(response.status, parseRetryAfter(response.headers.get("retry-after"))); @@ -206,25 +237,15 @@ export class KiconnectProvider implements LLMProvider, StructuredLLMProvider { throw new TutorProviderError("model-unavailable"); } return [...new Set(parsed.data.data.map(({ id }) => id))]; - } catch (error) { - if (error instanceof TutorProviderError) throw error; - if (error instanceof DOMException && error.name === "AbortError") { - throw new TutorProviderError("provider-timeout"); - } - if (error instanceof Error && error.name === "AbortError") { - throw new TutorProviderError("provider-timeout"); - } - throw new TutorProviderError("provider-unavailable"); - } finally { - globalThis.clearTimeout(timeout); - } + }); } async generateLearningQuestion( request: LLMProviderRequest, credential: string, + signal?: AbortSignal, ): Promise { - const model = await this.resolveModel(request.model, credential); + const model = await this.resolveModel(request.model, credential, signal); const body = await this.requestChatCompletion({ model, temperature: TUTOR_TEMPERATURE, @@ -232,13 +253,14 @@ export class KiconnectProvider implements LLMProvider, StructuredLLMProvider { { role: "system", content: request.systemPrompt }, { role: "user", content: request.userPrompt }, ], - }, credential); + }, credential, signal); return parseProviderQuestion(body, this.logger); } async generateStructuredResponse( request: StructuredLLMProviderRequest, credential: string, + signal?: AbortSignal, ): Promise { const body = await this.requestChatCompletion({ model: request.model, @@ -248,7 +270,7 @@ export class KiconnectProvider implements LLMProvider, StructuredLLMProvider { { role: "user", content: request.userPrompt }, ], response_format: { type: "json_object" }, - }, credential); + }, credential, signal); const parsed = structuredCompletionSchema.safeParse(body); if (!parsed.success) throw new TutorProviderError("invalid-response"); @@ -269,11 +291,9 @@ export class KiconnectProvider implements LLMProvider, StructuredLLMProvider { private async requestChatCompletion( requestBody: ChatCompletionRequestBody, credential: string, + signal?: AbortSignal, ): Promise { - const controller = new AbortController(); - const timeout = globalThis.setTimeout(() => controller.abort(), config.tutor.timeoutMs); - - try { + return withProviderTimeout(signal, async (providerSignal) => { const response = await fetch(`${this.baseUrl}/chat/completions`, { method: "POST", headers: { @@ -281,29 +301,18 @@ export class KiconnectProvider implements LLMProvider, StructuredLLMProvider { Authorization: `Bearer ${credential}`, }, body: JSON.stringify(requestBody), - signal: controller.signal, + signal: providerSignal, }); if (!response.ok) { throw providerErrorForStatus(response.status, parseRetryAfter(response.headers.get("retry-after"))); } return await response.json().catch(() => null); - } catch (error) { - if (error instanceof TutorProviderError) throw error; - if (error instanceof DOMException && error.name === "AbortError") { - throw new TutorProviderError("provider-timeout"); - } - if (error instanceof Error && error.name === "AbortError") { - throw new TutorProviderError("provider-timeout"); - } - throw new TutorProviderError("provider-unavailable"); - } finally { - globalThis.clearTimeout(timeout); - } + }); } - private async resolveModel(model: string, credential: string): Promise { + private async resolveModel(model: string, credential: string, signal?: AbortSignal): Promise { if (model !== "auto") return model; - const modelId = rankTutorModels(await this.listModels(credential))[0]; + const modelId = rankTutorModels(await this.listModels(credential, signal))[0]; if (!modelId) throw new TutorProviderError("model-unavailable"); return modelId; } diff --git a/server/services/tutor/llm-provider.ts b/server/services/tutor/llm-provider.ts index 3fb9852a9..6d559b3b2 100644 --- a/server/services/tutor/llm-provider.ts +++ b/server/services/tutor/llm-provider.ts @@ -27,6 +27,7 @@ export interface StructuredLLMProvider { generateStructuredResponse( request: StructuredLLMProviderRequest, credential: string, + signal?: AbortSignal, ): Promise; } @@ -49,9 +50,10 @@ export class TutorProviderError extends Error { } export interface LLMProvider { - listModels(credential: string): Promise; + listModels(credential: string, signal?: AbortSignal): Promise; generateLearningQuestion( request: LLMProviderRequest, credential: string, + signal?: AbortSignal, ): Promise; } diff --git a/server/services/tutor/tutor-service.ts b/server/services/tutor/tutor-service.ts index d39125103..fd53667d8 100644 --- a/server/services/tutor/tutor-service.ts +++ b/server/services/tutor/tutor-service.ts @@ -63,8 +63,13 @@ type TutorDialogArguments = [ requestedModel: string | undefined, difficulty?: TutorDifficulty, courseContent?: TutorPlanningContentContext, + signal?: AbortSignal, ]; +function throwIfTutorRequestAborted(signal: AbortSignal | undefined): void { + signal?.throwIfAborted(); +} + type TutorPromptOptions = { readonly didacticBrief?: TutorPlan; readonly strategy?: EffectiveTutorStrategy; @@ -655,11 +660,19 @@ export class TutorService { private readonly planningExtension?: TutorPlanningExtension, ) {} - private inSessionOrder(courseContent: TutorPlanningContentContext | undefined, request: () => Promise): Promise { + private inSessionOrder( + courseContent: TutorPlanningContentContext | undefined, + request: () => Promise, + signal?: AbortSignal, + ): Promise { const state = courseContent?.progressionState; - if (!state) return request(); + const run = () => { + throwIfTutorRequestAborted(signal); + return request(); + }; + if (!state) return run(); const previous = this.sessionQueues.get(state) ?? Promise.resolve(); - const current = previous.then(request, request); + const current = previous.then(run, run); this.sessionQueues.set(state, current.catch(() => undefined)); return current; } @@ -670,8 +683,13 @@ export class TutorService { requestedModel: string | undefined, difficulty: TutorDifficulty = TUTOR_DEFAULT_DIFFICULTY, courseContent?: TutorPlanningContentContext, + signal?: AbortSignal, ): Promise { - return this.inSessionOrder(courseContent, () => this.generateQuestionInSession(code, credential, requestedModel, difficulty, courseContent)); + return this.inSessionOrder( + courseContent, + () => this.generateQuestionInSession(code, credential, requestedModel, difficulty, courseContent, signal), + signal, + ); } private async generateQuestionInSession( @@ -680,17 +698,23 @@ export class TutorService { requestedModel: string | undefined, difficulty: TutorDifficulty, courseContent?: TutorPlanningContentContext, + signal?: AbortSignal, ): Promise { + throwIfTutorRequestAborted(signal); const requestCredential = this.resolveCredential(credential); const context = buildTutorContext(code); const transaction = beginTutorPlanningTransaction(courseContent); const planningResult = this.planningExtension ? await this.planningExtension.planInitial({ code, history: [], difficulty, courseContent: transaction.courseContent }) : null; - const strategy = strategyFromPlanningResult(planningResult) ?? await this.resolveStrategy(code, transaction.courseContent); + throwIfTutorRequestAborted(signal); + const strategy = strategyFromPlanningResult(planningResult) ?? await this.resolveStrategy(code, transaction.courseContent, signal); + throwIfTutorRequestAborted(signal); + const model = await this.resolveModel(requestedModel, requestCredential, signal); + throwIfTutorRequestAborted(signal); const providerResult: ProviderQuestionResult = await this.provider.generateLearningQuestion( { - model: await this.resolveModel(requestedModel, requestCredential), + model, systemPrompt: TUTOR_SYSTEM_PROMPT, userPrompt: buildUserPrompt( code, @@ -702,11 +726,14 @@ export class TutorService { ), }, requestCredential, + ...(signal ? [signal] : []), ); + throwIfTutorRequestAborted(signal); const validatedResult = stripProviderPlanningMetadata(validateLearningQuestion(providerResult.result, difficulty)); const plannedResult = planningResult ? applyPlanningOutcome(validatedResult, planningResult) : applyStrategyMetadata(validatedResult, strategy); + throwIfTutorRequestAborted(signal); transaction.commit(); const { answerRating: _initialAnswerRating, ...initialResult } = plannedResult; return { @@ -717,16 +744,18 @@ export class TutorService { } async generateDialogResponse(...args: TutorDialogArguments): Promise { - return this.inSessionOrder(args[7], () => this.generateDialogResponseInSession(...args)); + return this.inSessionOrder(args[7], () => this.generateDialogResponseInSession(...args), args[8]); } private async generateDialogResponseInSession(...args: TutorDialogArguments): Promise { - const [code, history, question, answer, credential, requestedModel, difficulty = TUTOR_DEFAULT_DIFFICULTY, courseContent] = args; + const [code, history, question, answer, credential, requestedModel, difficulty = TUTOR_DEFAULT_DIFFICULTY, courseContent, signal] = args; + throwIfTutorRequestAborted(signal); const requestCredential = this.resolveCredential(credential); const parsedHistory = history.map((entry) => tutorDialogTurnSchema.parse(entry)); const transaction = beginTutorPlanningTransaction(courseContent); if (isClearlyNonLearningAnswer(answer)) { - const strategy = await this.resolveStrategy(code, transaction.courseContent); + const strategy = await this.resolveStrategy(code, transaction.courseContent, signal); + throwIfTutorRequestAborted(signal); return { model: requestedModel ?? "fallback", result: applyStrategyMetadata(buildPhilosophicalFallback(parsedHistory, difficulty), strategy), @@ -736,11 +765,15 @@ export class TutorService { const context = buildTutorContext(code); // The provider rates the answer to `question`, so its didactic context must describe that // answered question, not a question planned afterwards. - const currentPlanningResult = await this.planAnsweredQuestion(code, parsedHistory, question, difficulty, transaction.courseContent); - const strategy = strategyFromPlanningResult(currentPlanningResult) ?? await this.resolveStrategy(code, transaction.courseContent); + const currentPlanningResult = await this.planAnsweredQuestion(code, parsedHistory, question, difficulty, transaction.courseContent, signal); + throwIfTutorRequestAborted(signal); + const strategy = strategyFromPlanningResult(currentPlanningResult) ?? await this.resolveStrategy(code, transaction.courseContent, signal); + throwIfTutorRequestAborted(signal); + const model = await this.resolveModel(requestedModel, requestCredential, signal); + throwIfTutorRequestAborted(signal); const providerResult = await this.provider.generateLearningQuestion( { - model: await this.resolveModel(requestedModel, requestCredential), + model, systemPrompt: TUTOR_SYSTEM_PROMPT, userPrompt: buildDialogPrompt(code, context, parsedHistory, question, answer, difficulty, { didacticBrief: currentPlanningResult && isTutorPlan(currentPlanningResult) ? currentPlanningResult : undefined, @@ -749,7 +782,9 @@ export class TutorService { }), }, requestCredential, + ...(signal ? [signal] : []), ); + throwIfTutorRequestAborted(signal); const validatedResult = stripProviderPlanningMetadata(validateLearningQuestion(providerResult.result, difficulty)); if (validatedResult.responseStyle === "normal" && validatedResult.answerRating === undefined) { throw new TutorProviderError("invalid-response"); @@ -761,8 +796,11 @@ export class TutorService { difficulty, strategy, courseContent: transaction.courseContent, + signal, }); + throwIfTutorRequestAborted(signal); const finalResult = followUp.result.strategyId ? followUp.result : applyStrategyMetadata(followUp.result, strategy); + throwIfTutorRequestAborted(signal); transaction.commit(); return { model: providerResult.model, @@ -782,6 +820,7 @@ export class TutorService { readonly difficulty: TutorDifficulty; readonly strategy: StrategyResolution; readonly courseContent?: TutorPlanningContentContext; + readonly signal?: AbortSignal; }, ): Promise<{ readonly result: TutorContentResult; readonly followUpSource: TutorFollowUpSource }> { if (validatedResult.responseStyle !== "normal") return { result: validatedResult, followUpSource: "provider" }; @@ -797,6 +836,7 @@ export class TutorService { difficulty: context.difficulty, courseContent: context.courseContent, }); + throwIfTutorRequestAborted(context.signal); if (nextPlan) { result = applyPlanningOutcome(result, nextPlan); if (isTutorPlan(nextPlan)) followUpSource = "planner"; @@ -811,16 +851,22 @@ export class TutorService { currentQuestion: string, difficulty: TutorDifficulty, courseContent?: TutorPlanningContentContext, + signal?: AbortSignal, ): Promise { // Never planInitial here: planning a new question while only evaluating an answer could // reserve content. Without planAnswered the prompt simply carries no didactic context. if (!this.planningExtension?.planAnswered) return null; - return this.planningExtension.planAnswered({ code, history, currentQuestion, difficulty, courseContent }); + const result = await this.planningExtension.planAnswered({ code, history, currentQuestion, difficulty, courseContent }); + throwIfTutorRequestAborted(signal); + return result; } - async getAvailableModels(credential: string | undefined): Promise { + async getAvailableModels(credential: string | undefined, signal?: AbortSignal): Promise { + throwIfTutorRequestAborted(signal); const requestCredential = this.resolveCredential(credential); - return this.provider.listModels(requestCredential); + const models = await this.provider.listModels(requestCredential, ...(signal ? [signal] : [])); + throwIfTutorRequestAborted(signal); + return models; } private resolveCredential(credential: string | undefined): string { @@ -831,18 +877,22 @@ export class TutorService { return requestCredential; } - private async resolveModel(requestedModel: string | undefined, credential: string): Promise { + private async resolveModel(requestedModel: string | undefined, credential: string, signal?: AbortSignal): Promise { const model = requestedModel ?? "auto"; if (model === "auto") return model; - const availableModels = await this.provider.listModels(credential); + const availableModels = await this.provider.listModels(credential, ...(signal ? [signal] : [])); + throwIfTutorRequestAborted(signal); return availableModels.includes(model) ? model : "auto"; } - private async resolveStrategy(code?: string, courseContent?: TutorPlanningContentContext): Promise { + private async resolveStrategy(code?: string, courseContent?: TutorPlanningContentContext, signal?: AbortSignal): Promise { if (this.planningExtension?.resolveStrategy) { try { - return await this.planningExtension.resolveStrategy({ code, courseContent }); + const strategy = await this.planningExtension.resolveStrategy({ code, courseContent }); + throwIfTutorRequestAborted(signal); + return strategy; } catch { + throwIfTutorRequestAborted(signal); // A strategy resolver is optional planning context; the built-in policy remains authoritative. } } diff --git a/tests/client/tutor.test.tsx b/tests/client/tutor.test.tsx index 933d38969..559e6b520 100644 --- a/tests/client/tutor.test.tsx +++ b/tests/client/tutor.test.tsx @@ -69,6 +69,38 @@ describe("useTutor", () => { ]); }); + it("aborts a pending dialog fetch when the hook unmounts", async () => { + let dialogSignal: AbortSignal | undefined; + const fetchMock = vi.spyOn(globalThis, "fetch").mockImplementation(async (input, init) => { + const url = String(input); + if (url === "/api/config") return new Response(JSON.stringify({ tutor: { provider: "kiconnect" } }), { status: 200 }); + if (url === "/api/tutor/models") return new Response(JSON.stringify({ models: ["pilot-model"] }), { status: 200 }); + if (url === "/api/tutor/question") return new Response(JSON.stringify({ + question: "Welche Ausgabe erwartest du?", + provider: "kiconnect", + model: "pilot-model", + }), { status: 200 }); + dialogSignal = init?.signal as AbortSignal | undefined; + return new Promise((_resolve, reject) => { + dialogSignal?.addEventListener("abort", () => reject(new DOMException("Aborted", "AbortError")), { once: true }); + }); + }); + const { result, unmount } = renderHook(() => useTutor()); + act(() => result.current.setCredential("personal-key")); + await act(async () => { + await result.current.loadModels(); + await result.current.generateQuestion("void setup(){} void loop(){}"); + }); + act(() => result.current.setAnswer("Antwort")); + act(() => { void result.current.submitAnswer("void setup(){} void loop(){}"); }); + await waitFor(() => expect(dialogSignal).toBeInstanceOf(AbortSignal)); + + unmount(); + + expect(dialogSignal?.aborted).toBe(true); + expect(fetchMock.mock.calls.filter(([input]) => String(input).startsWith("/api/tutor/")).every(([, init]) => init?.signal instanceof AbortSignal)).toBe(true); + }); + it("loads a valid persisted configured difficulty before starting a tutor session", () => { localStorage.setItem(TUTOR_CONFIGURED_DIFFICULTY_STORAGE_KEY, "72"); const { result } = renderHook(() => useTutor()); diff --git a/tests/server/examples/http-provider.test.ts b/tests/server/examples/http-provider.test.ts index 25d220627..557f4ec5f 100644 --- a/tests/server/examples/http-provider.test.ts +++ b/tests/server/examples/http-provider.test.ts @@ -52,4 +52,18 @@ describe("secure external examples fetcher", () => { await vi.advanceTimersByTimeAsync(5_001); await rejection; }); + + it("does not begin a Course Content fetch for an already-aborted request", async () => { + const controller = new AbortController(); + controller.abort(); + const fetchMock = vi.fn(async () => new Response("{}")); + vi.stubGlobal("fetch", fetchMock); + + await expect(new SecureExamplesFetcher().fetchText( + new URL("https://api.github.com/test"), + 100, + controller.signal, + )).rejects.toMatchObject({ name: "AbortError" }); + expect(fetchMock).not.toHaveBeenCalled(); + }); }); diff --git a/tests/server/examples/source-provider.test.ts b/tests/server/examples/source-provider.test.ts index 04abbeabb..54cbc931d 100644 --- a/tests/server/examples/source-provider.test.ts +++ b/tests/server/examples/source-provider.test.ts @@ -1,7 +1,7 @@ import { describe, expect, it, vi } from "vitest"; import { ExamplesCache } from "../../../server/services/examples/examples-cache"; import { ExamplesLoadController } from "../../../server/services/examples/examples-load-controller"; -import { SourceProvider } from "../../../server/services/examples/source-provider"; +import { SourceProvider, type RequestContext } from "../../../server/services/examples/source-provider"; const revisionA = "a".repeat(40); const revisionB = "b".repeat(40); @@ -18,7 +18,7 @@ function snapshot(marker: string) { } function harness( - resolve: (repository: string, ref: string) => Promise, + resolve: (repository: string, ref: string, context: RequestContext) => Promise, load: (_repository: string, revision: string) => Promise>, ) { let now = 0; @@ -87,6 +87,26 @@ describe("repository/ref source provider", () => { expect(loader).toHaveBeenCalledTimes(2); }); + it("does not mark a source stale when a request aborts during refresh", async () => { + const controller = new AbortController(); + const resolver = vi.fn() + .mockResolvedValueOnce(revisionA) + .mockImplementationOnce(async (_repository: string, _ref: string, request: RequestContext) => { + controller.abort(); + throw request.signal?.reason; + }) + .mockResolvedValue(revisionA); + const loader = vi.fn(async () => snapshot("a")); + const { provider, setNow } = harness(resolver, loader); + await provider.resolve("owner/repo", "main", context, false); + setNow(101); + + await expect(provider.resolve("owner/repo", "main", { ...context, signal: controller.signal }, false)) + .rejects.toMatchObject({ name: "AbortError" }); + await expect(provider.resolve("owner/repo", "main", context, false)).resolves.toMatchObject({ stale: false }); + expect(resolver).toHaveBeenCalledTimes(3); + }); + it("shares singleflight only for identical repository/ref sources", async () => { let release!: () => void; const gate = new Promise((resolve) => { release = resolve; }); diff --git a/tests/server/routes/tutor.routes.test.ts b/tests/server/routes/tutor.routes.test.ts index 5c77c8a32..33f9cc014 100644 --- a/tests/server/routes/tutor.routes.test.ts +++ b/tests/server/routes/tutor.routes.test.ts @@ -111,6 +111,7 @@ describe("Tutor HTTP route", () => { undefined, 30, undefined, + expect.any(AbortSignal), ); }); @@ -150,7 +151,7 @@ describe("Tutor HTTP route", () => { const response = await post(listening.url, "/api/tutor/models", { credential: "request-only-secret" }); expect(response).toEqual({ status: 200, body: { models: ["current-model", "another-model"] } }); - expect(service.getAvailableModels).toHaveBeenCalledWith("request-only-secret"); + expect(service.getAvailableModels).toHaveBeenCalledWith("request-only-secret", expect.any(AbortSignal)); }); it("requires a personal credential for model discovery", async () => { @@ -191,7 +192,7 @@ describe("Tutor HTTP route", () => { }); expect(response).toEqual({ status: 200, body: { models: ["pilot-model"] } }); - expect(service.getAvailableModels).toHaveBeenCalledWith("request-only-secret"); + expect(service.getAvailableModels).toHaveBeenCalledWith("request-only-secret", expect.any(AbortSignal)); }); it("accepts a personal credential over HTTP in the Docker test gateway bypass profile", async () => { @@ -206,7 +207,7 @@ describe("Tutor HTTP route", () => { }); expect(response).toEqual({ status: 200, body: { models: ["pilot-model"] } }); - expect(service.getAvailableModels).toHaveBeenCalledWith("request-only-secret"); + expect(service.getAvailableModels).toHaveBeenCalledWith("request-only-secret", expect.any(AbortSignal)); } finally { config.dockerTestBypassGateway = originalBypass; } @@ -250,6 +251,7 @@ describe("Tutor HTTP route", () => { undefined, 30, undefined, + expect.any(AbortSignal), ); }); @@ -299,7 +301,7 @@ describe("Tutor HTTP route", () => { const session = (first.body as { courseContentSession: string }).courseContentSession; expect(session).toMatch(/^[0-9a-f-]{36}$/); expect(service.generateQuestion).toHaveBeenCalledWith( - "void setup(){}", "request-only-secret", undefined, 30, contentA, + "void setup(){}", "request-only-secret", undefined, 30, contentA, expect.any(AbortSignal), ); const second = await post(listening.url, "/api/tutor/dialog", { @@ -313,7 +315,7 @@ describe("Tutor HTTP route", () => { expect(second.status).toBe(200); expect(service.generateDialogResponse).toHaveBeenCalledWith( "void setup(){}", [], "Frage A", "Antwort", - "request-only-secret", undefined, 30, expect.objectContaining(contentA), + "request-only-secret", undefined, 30, expect.objectContaining(contentA), expect.any(AbortSignal), ); expect(resolver.resolveTutorContent).toHaveBeenCalledOnce(); }); @@ -331,4 +333,53 @@ describe("Tutor HTTP route", () => { expect(response.status).toBe(400); expect(service.generateQuestion).not.toHaveBeenCalled(); }); + + it("aborts the shared route and Course Content signal when the client disconnects", async () => { + const content = { + repository: "owner/repo" as const, + ref: "main" as const, + revision: "c".repeat(40), + tutor: { status: "absent" as const }, + }; + let resolveGeneration: ((value: unknown) => void) | undefined; + const service = { + generateQuestion: vi.fn(( + _code: string, + _credential: string, + _model: string | undefined, + _difficulty: number, + _content: unknown, + _signal?: AbortSignal, + ) => new Promise((resolve) => { resolveGeneration = resolve; })), + }; + const resolver = { resolveTutorContent: vi.fn().mockResolvedValue(content) }; + const listening = await start(service, false, resolver); + server = listening.server; + const target = new URL("/api/tutor/question", listening.url); + const payload = JSON.stringify({ + code: "void setup(){}", + credential: "request-only-secret", + courseContent: { repository: "owner/repo", ref: "main", revision: "c".repeat(40) }, + }); + const request = http.request({ + hostname: target.hostname, + port: target.port, + path: target.pathname, + method: "POST", + headers: { "content-type": "application/json", "content-length": Buffer.byteLength(payload) }, + }, () => undefined); + request.on("error", () => undefined); + request.end(payload); + + try { + await vi.waitFor(() => expect(service.generateQuestion).toHaveBeenCalledOnce()); + const signal = service.generateQuestion.mock.calls[0]?.[5]; + request.destroy(); + await vi.waitFor(() => expect(signal?.aborted).toBe(true)); + expect(resolver.resolveTutorContent.mock.calls[0]?.[1].signal).toBe(signal); + } finally { + request.destroy(); + resolveGeneration?.({ model: "pilot-model", result: { question: "Frage" } }); + } + }); }); diff --git a/tests/server/services/course-content/course-content-loader.test.ts b/tests/server/services/course-content/course-content-loader.test.ts index 85f2c9e4b..30d0859c1 100644 --- a/tests/server/services/course-content/course-content-loader.test.ts +++ b/tests/server/services/course-content/course-content-loader.test.ts @@ -32,6 +32,26 @@ function fetcherFor(files: Record) { } describe("unified Course Content loader", () => { + it("rethrows an abort while loading optional Tutor content", async () => { + const controller = new AbortController(); + const fetchText = vi.fn(async (url: URL, _maxBytes: number, signal?: AbortSignal) => { + if (url.pathname.endsWith("/manifest.json")) { + return JSON.stringify({ + schemaVersion: 2, + examples: [example], + tutor: { manifest: "tutor/manifest.yaml" }, + }); + } + if (url.pathname.endsWith("/examples/main.ino")) return "void setup() {}"; + expect(signal).toBe(controller.signal); + controller.abort(); + throw controller.signal.reason; + }); + + await expect(new CourseContentLoader({ fetchText }, 2).load("owner/repo", revision, controller.signal)) + .rejects.toMatchObject({ name: "AbortError" }); + }); + it("loads schema-v1 Examples and reports no Tutor capability", async () => { const { fetchText } = fetcherFor({ [`/owner/repo/${revision}/manifest.json`]: JSON.stringify({ schemaVersion: 1, examples: [example] }), diff --git a/tests/server/services/tutor/kiconnect-provider.test.ts b/tests/server/services/tutor/kiconnect-provider.test.ts index dfa7005be..27de27083 100644 --- a/tests/server/services/tutor/kiconnect-provider.test.ts +++ b/tests/server/services/tutor/kiconnect-provider.test.ts @@ -24,6 +24,54 @@ describe("KiconnectProvider", () => { ); }); + it("composes caller cancellation with the model-list timeout", async () => { + const caller = new AbortController(); + const fetchMock = vi.fn((_url: string, init: RequestInit) => new Promise((_resolve, reject) => { + const signal = init.signal; + const timeout = setTimeout(() => reject(new Error("caller cancellation was not forwarded")), 250); + signal?.addEventListener("abort", () => { + clearTimeout(timeout); + reject(signal.reason); + }, { once: true }); + })); + vi.stubGlobal("fetch", fetchMock); + + const pending = new KiconnectProvider().listModels("request-key", caller.signal); + await vi.waitFor(() => expect(fetchMock).toHaveBeenCalledOnce()); + const providerSignal = (fetchMock.mock.calls[0] as [string, RequestInit])[1].signal; + caller.abort(); + + await expect(pending).rejects.toMatchObject({ name: "AbortError" }); + expect(providerSignal).not.toBe(caller.signal); + expect(providerSignal?.aborted).toBe(true); + }); + + it("composes caller cancellation with the completion timeout", async () => { + const caller = new AbortController(); + const fetchMock = vi.fn((_url: string, init: RequestInit) => new Promise((_resolve, reject) => { + const signal = init.signal; + const timeout = setTimeout(() => reject(new Error("caller cancellation was not forwarded")), 250); + signal?.addEventListener("abort", () => { + clearTimeout(timeout); + reject(signal.reason); + }, { once: true }); + })); + vi.stubGlobal("fetch", fetchMock); + + const pending = new KiconnectProvider().generateLearningQuestion({ + model: "pilot-model", + systemPrompt: "system", + userPrompt: "user", + }, "request-key", caller.signal); + await vi.waitFor(() => expect(fetchMock).toHaveBeenCalledOnce()); + const providerSignal = (fetchMock.mock.calls[0] as [string, RequestInit])[1].signal; + caller.abort(); + + await expect(pending).rejects.toMatchObject({ name: "AbortError" }); + expect(providerSignal).not.toBe(caller.signal); + expect(providerSignal?.aborted).toBe(true); + }); + it("uses the server-configured OpenAI-compatible endpoint and parses structured JSON", async () => { const fetchMock = vi.fn() .mockResolvedValueOnce(new Response(JSON.stringify({ data: [{ id: "pilot-model" }] }), { status: 200 })) diff --git a/tests/server/services/tutor/session-concurrency.test.ts b/tests/server/services/tutor/session-concurrency.test.ts index 453957add..9e952f9c0 100644 --- a/tests/server/services/tutor/session-concurrency.test.ts +++ b/tests/server/services/tutor/session-concurrency.test.ts @@ -37,6 +37,63 @@ function delayedProvider(delayByQuestion: Record): LLMProvider { } describe("Tutor session concurrency", () => { + it("does not commit progression when the request aborts before the provider result is consumed", async () => { + const { context } = await pinnedSession(); + const before = structuredClone(context.progressionState); + const controller = new AbortController(); + let notifyProviderStarted!: (signal: AbortSignal | undefined) => void; + let releaseProvider!: () => void; + const providerStarted = new Promise((resolve) => { notifyProviderStarted = resolve; }); + const provider: LLMProvider = { + listModels: vi.fn().mockResolvedValue(["pilot-model"]), + generateLearningQuestion: vi.fn(async (_request: LLMProviderRequest, _credential: string, signal?: AbortSignal) => { + notifyProviderStarted(signal); + await new Promise((resolve) => { releaseProvider = resolve; }); + return { model: "pilot-model", result: { question: "Welche Folge erwartest du?", responseStyle: "normal" as const } }; + }), + }; + const service = new TutorService(provider, new CurriculumTutorAdapter()); + const pending = service.generateQuestion(CODE, "key", undefined, 30, context, controller.signal); + + const providerSignal = await providerStarted; + controller.abort(); + releaseProvider(); + + await expect(pending).rejects.toMatchObject({ name: "AbortError" }); + expect(providerSignal).toBe(controller.signal); + expect(context.progressionState).toEqual(before); + }); + + it("does not start a same-session queued request after it has been aborted", async () => { + const { context } = await pinnedSession(); + let releaseFirst!: () => void; + let notifyFirstStarted!: () => void; + let providerCalls = 0; + const firstStarted = new Promise((resolve) => { notifyFirstStarted = resolve; }); + const provider: LLMProvider = { + listModels: vi.fn().mockResolvedValue(["pilot-model"]), + generateLearningQuestion: vi.fn(async (_request: LLMProviderRequest, _credential: string, _signal?: AbortSignal) => { + providerCalls += 1; + if (providerCalls === 1) { + notifyFirstStarted(); + await new Promise((resolve) => { releaseFirst = resolve; }); + } + return { model: "pilot-model", result: { question: "Welche Folge erwartest du?", responseStyle: "normal" as const } }; + }), + }; + const service = new TutorService(provider); + const first = service.generateQuestion(CODE, "key", undefined, 30, context); + await firstStarted; + const controller = new AbortController(); + const queued = service.generateQuestion(CODE, "key", undefined, 30, context, controller.signal); + controller.abort(); + releaseFirst(); + + await first; + await expect(queued).rejects.toMatchObject({ name: "AbortError" }); + expect(providerCalls).toBe(1); + }); + it("keeps the evidence of two answers that overlap on the same session", async () => { const { context, questions } = await pinnedSession(); const first = questions["int-width-direct"]!; From 3aaf14fc5639c2f20d5f98a4416a3e5a19382bf5 Mon Sep 17 00:00:00 2001 From: ttbombadil Date: Sun, 4 Oct 2026 16:13:56 +0200 Subject: [PATCH 2/5] fix(tutor): preserve shared loads across cancellation --- server/services/examples/examples-cache.ts | 127 +++++++++++++++--- server/services/examples/source-provider.ts | 16 ++- tests/server/examples/source-provider.test.ts | 49 +++++++ 3 files changed, 170 insertions(+), 22 deletions(-) diff --git a/server/services/examples/examples-cache.ts b/server/services/examples/examples-cache.ts index 3266d06eb..9c7a0d408 100644 --- a/server/services/examples/examples-cache.ts +++ b/server/services/examples/examples-cache.ts @@ -43,11 +43,18 @@ export interface ExamplesCacheOptions { now?: () => number; } +interface SharedFlight { + controller: AbortController; + promise: Promise; + subscribers: number; + settled: boolean; +} + export class ExamplesCache { private readonly sources = new Map(); private readonly revisions = new Map(); - private readonly sourceFlights = new Map>(); - private readonly revisionFlights = new Map>(); + private readonly sourceFlights = new Map(); + private readonly revisionFlights = new Map(); private totalSnapshotBytes = 0; private readonly now: () => number; @@ -135,21 +142,23 @@ export class ExamplesCache { return stored; } - async withSourceSingleflight(key: SourceCacheKey, load: () => Promise): Promise { - const current = this.sourceFlights.get(key) as Promise | undefined; - if (current) return current; - this.admitSource(key); - const flight = load().finally(() => this.sourceFlights.delete(key)); - this.sourceFlights.set(key, flight); - return flight; + async withSourceSingleflight( + key: SourceCacheKey, + load: (signal: AbortSignal) => Promise, + requestSignal?: AbortSignal, + ): Promise { + if (requestSignal?.aborted) return Promise.reject(requestSignal.reason ?? new DOMException("Request aborted", "AbortError")); + const current = this.sourceFlights.get(key); + if (!current || current.controller.signal.aborted) this.admitSource(key); + return this.withSingleflight(this.sourceFlights, key, load, requestSignal); } - async withRevisionSingleflight(key: RevisionCacheKey, load: () => Promise): Promise { - const current = this.revisionFlights.get(key) as Promise | undefined; - if (current) return current; - const flight = load().finally(() => this.revisionFlights.delete(key)); - this.revisionFlights.set(key, flight); - return flight; + async withRevisionSingleflight( + key: RevisionCacheKey, + load: (signal: AbortSignal) => Promise, + requestSignal?: AbortSignal, + ): Promise { + return this.withSingleflight(this.revisionFlights, key, load, requestSignal); } getStats() { @@ -172,6 +181,94 @@ export class ExamplesCache { } } + private withSingleflight( + flights: Map, + key: K, + load: (signal: AbortSignal) => Promise, + requestSignal?: AbortSignal, + ): Promise { + if (requestSignal?.aborted) return Promise.reject(requestSignal.reason ?? new DOMException("Request aborted", "AbortError")); + + let flight = flights.get(key) as SharedFlight & { promise: Promise } | undefined; + if (flight?.controller.signal.aborted) { + if (flights.get(key) === flight) flights.delete(key); + flight = undefined; + } + if (!flight) { + const controller = new AbortController(); + let created!: SharedFlight & { promise: Promise }; + let resolveFlight!: (value: T) => void; + let rejectFlight!: (error: unknown) => void; + const promise = new Promise((resolve, reject) => { + resolveFlight = resolve; + rejectFlight = reject; + }); + const settle = () => { + created.settled = true; + if (flights.get(key) === created) flights.delete(key); + }; + created = { controller, promise, subscribers: 0, settled: false }; + flight = created; + flights.set(key, created); + try { + void Promise.resolve(load(controller.signal)).then( + (value) => { + settle(); + resolveFlight(value); + }, + (error: unknown) => { + settle(); + rejectFlight(error); + }, + ); + } catch (error) { + settle(); + rejectFlight(error); + } + } + + return this.subscribeToFlight(flight, requestSignal); + } + + private subscribeToFlight( + flight: SharedFlight & { promise: Promise }, + requestSignal?: AbortSignal, + ): Promise { + if (requestSignal?.aborted) return Promise.reject(requestSignal.reason ?? new DOMException("Request aborted", "AbortError")); + flight.subscribers += 1; + + return new Promise((resolve, reject) => { + let left = false; + const leave = () => { + if (left) return; + left = true; + requestSignal?.removeEventListener("abort", onAbort); + flight.subscribers -= 1; + if (flight.subscribers === 0 && !flight.settled) flight.controller.abort(); + }; + const onAbort = () => { + leave(); + reject(requestSignal?.reason ?? new DOMException("Request aborted", "AbortError")); + }; + + requestSignal?.addEventListener("abort", onAbort, { once: true }); + if (requestSignal?.aborted) { + onAbort(); + return; + } + flight.promise.then( + (value) => { + leave(); + resolve(value as T); + }, + (error: unknown) => { + leave(); + reject(error); + }, + ); + }); + } + private enforceRevisionLimits(): void { while (this.revisions.size > this.options.maxSnapshots || this.totalSnapshotBytes > this.options.maxSnapshotBytes) { if (this.evictOldestUnpinnedRevision()) continue; diff --git a/server/services/examples/source-provider.ts b/server/services/examples/source-provider.ts index 722930eb7..743688abe 100644 --- a/server/services/examples/source-provider.ts +++ b/server/services/examples/source-provider.ts @@ -69,12 +69,13 @@ export class SourceProvider { } try { - return await this.cache.withSourceSingleflight(sourceKey, () => + return await this.cache.withSourceSingleflight(sourceKey, (signal) => this.loads.runLoad(async () => { - const revision = await this.resolver.resolve(repository, ref, context); + const loadContext = { ...context, signal }; + const revision = await this.resolver.resolve(repository, ref, loadContext); const revisionKey = toRevisionCacheKey(repository, revision); const cachedSnapshot = this.cache.getRevision(revisionKey); - const loadedSnapshot = cachedSnapshot ?? await this.loadRevision(repository, revision, context); + const loadedSnapshot = cachedSnapshot ?? await this.loadRevision(repository, revision, loadContext); const checkedAt = this.now(); const activated = this.cache.activateSource(sourceKey, { repository, @@ -91,7 +92,8 @@ export class SourceProvider { status: cachedSnapshot ? "cache" as const : "remote" as const, stale: false, }; - }, context.signal), + }, signal), + context.signal, ); } catch (error) { if (context.signal?.aborted) throw context.signal.reason ?? error; @@ -126,10 +128,10 @@ export class SourceProvider { context: RequestContext, ): Promise> { const key = toRevisionCacheKey(repository, revision); - return this.cache.withRevisionSingleflight(key, async () => { - const loaded = await this.revisionLoader.load(repository, revision, context.signal); + return this.cache.withRevisionSingleflight(key, async (signal) => { + const loaded = await this.revisionLoader.load(repository, revision, signal); return { repository, revision, ...loaded }; - }); + }, context.signal); } private fromEntry(entry: SourceCacheEntry): ResolvedSourceSnapshot { diff --git a/tests/server/examples/source-provider.test.ts b/tests/server/examples/source-provider.test.ts index 54cbc931d..cd5d7f645 100644 --- a/tests/server/examples/source-provider.test.ts +++ b/tests/server/examples/source-provider.test.ts @@ -107,6 +107,55 @@ describe("repository/ref source provider", () => { expect(resolver).toHaveBeenCalledTimes(3); }); + it("keeps a shared refresh alive while another request still needs it", async () => { + let release!: () => void; + let loadSignal!: AbortSignal; + const gate = new Promise((resolve) => { release = resolve; }); + const resolver = vi.fn(async (_repository: string, _ref: string, request: RequestContext) => { + loadSignal = request.signal!; + await gate; + return revisionA; + }); + const loader = vi.fn(async () => snapshot("a")); + const { provider } = harness(resolver, loader); + const firstController = new AbortController(); + const secondController = new AbortController(); + const first = provider.resolve("owner/repo", "main", { ...context, signal: firstController.signal }, false); + const second = provider.resolve("owner/repo", "main", { ...context, requestId: "request-b", signal: secondController.signal }, false); + await vi.waitFor(() => expect(resolver).toHaveBeenCalledTimes(1)); + + firstController.abort(new DOMException("First client left", "AbortError")); + expect(loadSignal.aborted).toBe(false); + release(); + + await expect(first).rejects.toMatchObject({ name: "AbortError" }); + await expect(second).resolves.toMatchObject({ revision: revisionA, stale: false }); + expect(loader).toHaveBeenCalledTimes(1); + }); + + it("aborts a shared refresh after its final request leaves", async () => { + let loadSignal!: AbortSignal; + const resolver = vi.fn(async (_repository: string, _ref: string, request: RequestContext) => { + loadSignal = request.signal!; + return new Promise((_resolve, reject) => { + loadSignal.addEventListener("abort", () => reject(loadSignal.reason), { once: true }); + }); + }); + const { provider } = harness(resolver, vi.fn(async () => snapshot("a"))); + const firstController = new AbortController(); + const secondController = new AbortController(); + const first = provider.resolve("owner/repo", "main", { ...context, signal: firstController.signal }, false); + const second = provider.resolve("owner/repo", "main", { ...context, requestId: "request-b", signal: secondController.signal }, false); + await vi.waitFor(() => expect(resolver).toHaveBeenCalledTimes(1)); + + firstController.abort(new DOMException("First client left", "AbortError")); + secondController.abort(new DOMException("Second client left", "AbortError")); + + await expect(first).rejects.toMatchObject({ name: "AbortError" }); + await expect(second).rejects.toMatchObject({ name: "AbortError" }); + expect(loadSignal.aborted).toBe(true); + }); + it("shares singleflight only for identical repository/ref sources", async () => { let release!: () => void; const gate = new Promise((resolve) => { release = resolve; }); From 204b111a9c080ff54511420b65807d15ef6c12b1 Mon Sep 17 00:00:00 2001 From: ttbombadil Date: Sun, 4 Oct 2026 16:16:55 +0200 Subject: [PATCH 3/5] fix(tutor): prevent canceled refresh commits --- server/services/examples/source-provider.ts | 2 ++ tests/server/examples/source-provider.test.ts | 26 ++++++++++++++++++- 2 files changed, 27 insertions(+), 1 deletion(-) diff --git a/server/services/examples/source-provider.ts b/server/services/examples/source-provider.ts index 743688abe..287f844ec 100644 --- a/server/services/examples/source-provider.ts +++ b/server/services/examples/source-provider.ts @@ -73,9 +73,11 @@ export class SourceProvider { this.loads.runLoad(async () => { const loadContext = { ...context, signal }; const revision = await this.resolver.resolve(repository, ref, loadContext); + signal.throwIfAborted(); const revisionKey = toRevisionCacheKey(repository, revision); const cachedSnapshot = this.cache.getRevision(revisionKey); const loadedSnapshot = cachedSnapshot ?? await this.loadRevision(repository, revision, loadContext); + signal.throwIfAborted(); const checkedAt = this.now(); const activated = this.cache.activateSource(sourceKey, { repository, diff --git a/tests/server/examples/source-provider.test.ts b/tests/server/examples/source-provider.test.ts index cd5d7f645..6967f8ca5 100644 --- a/tests/server/examples/source-provider.test.ts +++ b/tests/server/examples/source-provider.test.ts @@ -1,5 +1,5 @@ import { describe, expect, it, vi } from "vitest"; -import { ExamplesCache } from "../../../server/services/examples/examples-cache"; +import { ExamplesCache, toSourceCacheKey } from "../../../server/services/examples/examples-cache"; import { ExamplesLoadController } from "../../../server/services/examples/examples-load-controller"; import { SourceProvider, type RequestContext } from "../../../server/services/examples/source-provider"; @@ -156,6 +156,30 @@ describe("repository/ref source provider", () => { expect(loadSignal.aborted).toBe(true); }); + it("does not activate a source after every request cancels", async () => { + let release!: () => void; + const gate = new Promise((resolve) => { release = resolve; }); + const resolver = vi.fn() + .mockResolvedValueOnce(revisionA) + .mockImplementationOnce(async () => { + await gate; + return revisionA; + }); + const { provider, cache, setNow } = harness(resolver, vi.fn(async () => snapshot("a"))); + await provider.resolve("owner/repo", "main", context, false); + setNow(101); + const controller = new AbortController(); + const pending = provider.resolve("owner/repo", "main", { ...context, signal: controller.signal }, false); + await vi.waitFor(() => expect(resolver).toHaveBeenCalledTimes(1)); + + controller.abort(new DOMException("Client left", "AbortError")); + await expect(pending).rejects.toMatchObject({ name: "AbortError" }); + release(); + await vi.waitFor(() => expect(cache.getStats().sourceFlights).toBe(0)); + + expect(cache.getSource(toSourceCacheKey("owner/repo", "main"))?.checkedAt).toBe(0); + }); + it("shares singleflight only for identical repository/ref sources", async () => { let release!: () => void; const gate = new Promise((resolve) => { release = resolve; }); From ffc967c7bf02c6ffa617d49808f12a28c89fb278 Mon Sep 17 00:00:00 2001 From: ttbombadil Date: Sun, 4 Oct 2026 16:21:45 +0200 Subject: [PATCH 4/5] fix(tutor): satisfy cancellation quality gate --- server/routes/tutor.routes.ts | 2 +- server/services/examples/examples-cache.ts | 11 ++++------- tests/server/examples/examples-cache.test.ts | 2 +- 3 files changed, 6 insertions(+), 9 deletions(-) diff --git a/server/routes/tutor.routes.ts b/server/routes/tutor.routes.ts index 14acab594..82b25b8ac 100644 --- a/server/routes/tutor.routes.ts +++ b/server/routes/tutor.routes.ts @@ -118,7 +118,7 @@ function bindRequestAbortSignal(req: Request, res: Response): AbortSignal { req.once("aborted", abort); res.once("close", onClose); res.once("finish", cleanup); - if (req.aborted || res.destroyed) abort(); + if ((req.destroyed && !req.complete) || res.destroyed) abort(); return controller.signal; } diff --git a/server/services/examples/examples-cache.ts b/server/services/examples/examples-cache.ts index 9c7a0d408..5736c02af 100644 --- a/server/services/examples/examples-cache.ts +++ b/server/services/examples/examples-cache.ts @@ -147,7 +147,7 @@ export class ExamplesCache { load: (signal: AbortSignal) => Promise, requestSignal?: AbortSignal, ): Promise { - if (requestSignal?.aborted) return Promise.reject(requestSignal.reason ?? new DOMException("Request aborted", "AbortError")); + if (requestSignal?.aborted) throw requestSignal.reason ?? new DOMException("Request aborted", "AbortError"); const current = this.sourceFlights.get(key); if (!current || current.controller.signal.aborted) this.admitSource(key); return this.withSingleflight(this.sourceFlights, key, load, requestSignal); @@ -210,8 +210,9 @@ export class ExamplesCache { created = { controller, promise, subscribers: 0, settled: false }; flight = created; flights.set(key, created); - try { - void Promise.resolve(load(controller.signal)).then( + void Promise.resolve() + .then(() => load(controller.signal)) + .then( (value) => { settle(); resolveFlight(value); @@ -221,10 +222,6 @@ export class ExamplesCache { rejectFlight(error); }, ); - } catch (error) { - settle(); - rejectFlight(error); - } } return this.subscribeToFlight(flight, requestSignal); diff --git a/tests/server/examples/examples-cache.test.ts b/tests/server/examples/examples-cache.test.ts index 7d25aa7bc..bc7787498 100644 --- a/tests/server/examples/examples-cache.test.ts +++ b/tests/server/examples/examples-cache.test.ts @@ -58,7 +58,7 @@ describe("external examples cache", () => { const first = cache.withSourceSingleflight(key, load); const second = cache.withSourceSingleflight(key, load); const different = cache.withSourceSingleflight(toSourceCacheKey("owner/repo", "preview"), load); - expect(loads).toBe(2); + await vi.waitFor(() => expect(loads).toBe(2)); release(); await Promise.all([first, second, different]); expect(loads).toBe(2); From f7a10eb7a547c790162336cb111ad50cab6325f4 Mon Sep 17 00:00:00 2001 From: ttbombadil Date: Sun, 4 Oct 2026 16:37:08 +0200 Subject: [PATCH 5/5] fix(tutor): distinguish request abort from fetch timeout --- .../course-content/course-content-loader.ts | 2 +- .../course-content-loader.test.ts | 24 +++++++++++++++++++ 2 files changed, 25 insertions(+), 1 deletion(-) diff --git a/server/services/course-content/course-content-loader.ts b/server/services/course-content/course-content-loader.ts index 5016d57ee..6584ee1c6 100644 --- a/server/services/course-content/course-content-loader.ts +++ b/server/services/course-content/course-content-loader.ts @@ -142,7 +142,7 @@ export class CourseContentLoader { validateTutorReferences(manifest, topicEntries, strategies, annotations); return { status: "valid", manifest, topics: topicEntries, strategies }; } catch (error) { - if (signal?.aborted || (error instanceof Error && error.name === "AbortError")) { + if (signal?.aborted) { throw signal?.reason ?? error; } return { status: "invalid", reason: "Tutor capability is invalid" }; diff --git a/tests/server/services/course-content/course-content-loader.test.ts b/tests/server/services/course-content/course-content-loader.test.ts index 30d0859c1..f97ffc01c 100644 --- a/tests/server/services/course-content/course-content-loader.test.ts +++ b/tests/server/services/course-content/course-content-loader.test.ts @@ -52,6 +52,30 @@ describe("unified Course Content loader", () => { .rejects.toMatchObject({ name: "AbortError" }); }); + it("keeps valid Examples when optional Tutor fetch times out with AbortError", async () => { + const controller = new AbortController(); + const fetchText = vi.fn(async (url: URL) => { + if (url.pathname.endsWith("/manifest.json")) { + return JSON.stringify({ + schemaVersion: 2, + examples: [example], + tutor: { manifest: "tutor/manifest.yaml" }, + }); + } + if (url.pathname.endsWith("/examples/main.ino")) return "void setup() {}"; + const timeout = new Error("request timed out"); + timeout.name = "AbortError"; + throw timeout; + }); + + const loaded = await new CourseContentLoader({ fetchText }, 2) + .load("owner/repo", revision, controller.signal); + + expect(controller.signal.aborted).toBe(false); + expect(loaded.examples).toHaveLength(1); + expect(loaded.tutor).toMatchObject({ status: "invalid" }); + }); + it("loads schema-v1 Examples and reports no Tutor capability", async () => { const { fetchText } = fetcherFor({ [`/owner/repo/${revision}/manifest.json`]: JSON.stringify({ schemaVersion: 1, examples: [example] }),