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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
52 changes: 40 additions & 12 deletions src/modules/quizzes/ai-client.ts
Original file line number Diff line number Diff line change
@@ -1,16 +1,20 @@
import { z } from "zod";
import { config } from "../../config/index.js";
import { logger } from "../../utils/logger.js";
import { createTransientRetryPolicy, createCircuitBreaker } from "../../utils/resilience.js";

interface AiQuizQuestion {
prompt: string;
options: string[];
correct_index: number;
}
const aiQuizQuestionSchema = z.object({
prompt: z.string(),
options: z.array(z.string()),
correct_index: z.number().int(),
});

interface AiQuizResponse {
quiz_id: string;
questions: AiQuizQuestion[];
}
const aiQuizResponseSchema = z.object({
quiz_id: z.string(),
questions: z.array(aiQuizQuestionSchema),
});

export type AiQuizQuestion = z.infer<typeof aiQuizQuestionSchema>;

export type AiDifficulty = "beginner" | "intermediate" | "advanced";

Expand All @@ -22,7 +26,10 @@ export interface GenerateQuizFromAIParams {
numQuestions: number;
}

export async function generateQuizFromAI(
const aiRetry = createTransientRetryPolicy("AI service", { maxAttempts: 3 });
const aiBreaker = createCircuitBreaker({ label: "AI service" });

async function requestQuiz(
params: GenerateQuizFromAIParams
): Promise<AiQuizQuestion[]> {
const controller = new AbortController();
Expand Down Expand Up @@ -50,8 +57,17 @@ export async function generateQuizFromAI(
throw new Error(`AI service returned ${response.status}`);
}

const data = (await response.json()) as AiQuizResponse;
return data.questions;
const raw = await response.json();
const parsed = aiQuizResponseSchema.safeParse(raw);
if (!parsed.success) {
logger.error(
{ issues: parsed.error.issues },
"AI service returned a malformed response"
);
throw new Error("AI service returned a malformed response");
}

return parsed.data.questions;
} catch (err) {
if (err instanceof Error && err.name === "AbortError") {
logger.error({ timeout: config.AI_TIMEOUT_MS }, "AI service request timed out");
Expand All @@ -62,3 +78,15 @@ export async function generateQuizFromAI(
clearTimeout(timeout);
}
}

/**
* Requests a generated quiz from the chainlearn-ai service. Wrapped with the
* same retry + circuit breaker pattern used for Stellar calls: transient
* failures are retried with backoff, and a persistent outage trips the
* breaker so requests fail fast instead of blocking on every quiz request.
*/
export async function generateQuizFromAI(
params: GenerateQuizFromAIParams
): Promise<AiQuizQuestion[]> {
return aiBreaker.execute(() => aiRetry.execute(() => requestQuiz(params)));
}
9 changes: 6 additions & 3 deletions src/services/retry-queue.ts
Original file line number Diff line number Diff line change
Expand Up @@ -50,15 +50,17 @@ export async function requeueReward(job: RetryJob): Promise<void> {

let processorRunning = false;
let processorTimer: ReturnType<typeof setTimeout> | null = null;
let processorGeneration = 0;

export async function startRetryProcessor(
processFn: (job: RetryJob) => Promise<boolean>
): Promise<void> {
if (processorRunning) return;
processorRunning = true;
const generation = ++processorGeneration;

const tick = async () => {
if (!processorRunning) return;
if (generation !== processorGeneration) return;
try {
const job = await dequeueReward();
if (job) {
Expand All @@ -70,17 +72,18 @@ export async function startRetryProcessor(
} catch (err) {
logger.error({ err }, "Retry processor tick failed");
}
if (processorRunning) {
if (generation === processorGeneration) {
processorTimer = setTimeout(tick, RETRY_INTERVAL_MS);
}
};

tick();
await tick();
logger.info("Retry processor started");
}

export function stopRetryProcessor(): void {
processorRunning = false;
processorGeneration++;
if (processorTimer) {
clearTimeout(processorTimer);
processorTimer = null;
Expand Down
149 changes: 16 additions & 133 deletions src/stellar/resilience.ts
Original file line number Diff line number Diff line change
@@ -1,144 +1,27 @@
import { retry, handleType, ExponentialBackoff } from "cockatiel";
import { logger } from "../utils/logger.js";
import {
createTransientRetryPolicy,
createCircuitBreaker,
withTimeout,
isCircuitBreakerError,
CircuitState,
CircuitBreakerOpenError,
TimeoutError,
} from "../utils/resilience.js";

function isTransientError(err: Error): boolean {
const name = err.name ?? "";
const msg = err.message ?? "";
export const stellarRetry = createTransientRetryPolicy("Stellar");

if (name === "FetchError" || name === "HttpError") return true;
if (
msg.includes("ECONNREFUSED") ||
msg.includes("ETIMEDOUT") ||
msg.includes("ECONNRESET") ||
msg.includes("ENOTFOUND") ||
msg.includes("socket hang up")
) {
return true;
}

const statusMatch = msg.match(/\b(502|503|504)\b/);
if (statusMatch) return true;

return false;
}

export const stellarRetry = retry(
handleType(Error, (err) => {
if (isTransientError(err)) {
logger.warn({ error: err.message }, "Stellar call retrying after transient error");
return true;
}
return false;
}),
{ backoff: new ExponentialBackoff() }
);

// Circuit breaker implementation
export enum CircuitState {
Closed = "Closed",
Open = "Open",
HalfOpen = "HalfOpen",
}

let circuitState = CircuitState.Closed;
let failureCount = 0;
let lastFailureTime = 0;
const THRESHOLD = 5;
const HALF_OPEN_AFTER = 30_000;

function recordSuccess(): void {
failureCount = 0;
if (circuitState !== CircuitState.Closed) {
logger.info("Circuit breaker reset to closed");
circuitState = CircuitState.Closed;
}
}

function recordFailure(): void {
failureCount++;
lastFailureTime = Date.now();

// Transition to Open from Closed state when threshold is reached
if (failureCount >= THRESHOLD && circuitState === CircuitState.Closed) {
circuitState = CircuitState.Open;
logger.warn("Circuit breaker opened after consecutive failures");
}

// If probe fails in HalfOpen state, re-open the circuit
if (circuitState === CircuitState.HalfOpen) {
circuitState = CircuitState.Open;
logger.warn("Circuit breaker re-opened after probe failure in HalfOpen state");
}
}
const stellarBreaker = createCircuitBreaker({ label: "Stellar" });

export function getCircuitState(): CircuitState {
if (circuitState === CircuitState.Open) {
if (Date.now() - lastFailureTime > HALF_OPEN_AFTER) {
circuitState = CircuitState.HalfOpen;
logger.info("Circuit breaker half-open — allowing probe request");
}
}
return circuitState;
return stellarBreaker.getState();
}

export function resetCircuitBreaker(): void {
circuitState = CircuitState.Closed;
failureCount = 0;
lastFailureTime = 0;
stellarBreaker.reset();
}

export async function circuitBreakerExecute<T>(fn: () => Promise<T>): Promise<T> {
const state = getCircuitState();

if (state === CircuitState.Open) {
throw new CircuitBreakerOpenError("Circuit breaker is open");
}

try {
const result = await fn();
recordSuccess();
return result;
} catch (err) {
if (err instanceof Error && isTransientError(err)) {
recordFailure();
}
throw err;
}
}

export function withTimeout<T>(promise: Promise<T>, ms: number): Promise<T> {
return new Promise<T>((resolve, reject) => {
const timer = setTimeout(() => {
reject(new TimeoutError(`Operation timed out after ${ms}ms`));
}, ms);

promise.then(
(value) => {
clearTimeout(timer);
resolve(value);
},
(err) => {
clearTimeout(timer);
reject(err);
},
);
});
export function circuitBreakerExecute<T>(fn: () => Promise<T>): Promise<T> {
return stellarBreaker.execute(fn);
}

export class CircuitBreakerOpenError extends Error {
constructor(message = "Circuit breaker is open") {
super(message);
this.name = "CircuitBreakerOpenError";
}
}

export class TimeoutError extends Error {
constructor(message = "Operation timed out") {
super(message);
this.name = "TimeoutError";
}
}

export function isCircuitBreakerError(err: unknown): boolean {
return err instanceof CircuitBreakerOpenError;
}
export { withTimeout, isCircuitBreakerError, CircuitState, CircuitBreakerOpenError, TimeoutError };
Loading