Skip to content
Open
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
20 changes: 17 additions & 3 deletions src/deeplake-api.ts
Original file line number Diff line number Diff line change
Expand Up @@ -174,6 +174,10 @@ const RETRYABLE_CODES = new Set([429, 500, 502, 503, 504]);
const MAX_RETRIES = 3;
const BASE_DELAY_MS = 500;
const MAX_CONCURRENCY = 5;
const DEFAULT_QUERY_TIMEOUT_MS = 10_000;
// Node's timers use a signed 32-bit delay. AbortSignal.timeout() accepts a
// larger unsigned range, but values above this overflow to a 1ms timer.
const MAX_TIMER_DELAY_MS = 2_147_483_647;

// Lazy read: the openclaw bundle replaces `process.env.HIVEMIND_QUERY_TIMEOUT_MS`
// with a `globalThis.__hivemind_tuning__?.HIVEMIND_QUERY_TIMEOUT_MS` lookup via
Expand All @@ -182,7 +186,12 @@ const MAX_CONCURRENCY = 5;
// Was previously `const QUERY_TIMEOUT_MS = …` at module top — that would have
// frozen the value to 10000 for the openclaw bundle regardless of pluginConfig.
function getQueryTimeoutMs(): number {
return Number(process.env.HIVEMIND_QUERY_TIMEOUT_MS ?? 10_000);
const configured = process.env.HIVEMIND_QUERY_TIMEOUT_MS;
if (configured === undefined || String(configured).trim() === "") return DEFAULT_QUERY_TIMEOUT_MS;

const parsed = Number(configured);
if (!Number.isFinite(parsed) || parsed < 0) return DEFAULT_QUERY_TIMEOUT_MS;
return Math.min(Math.trunc(parsed), MAX_TIMER_DELAY_MS);
Comment thread
coderabbitai[bot] marked this conversation as resolved.
}

function sleep(ms: number, signal?: AbortSignal): Promise<void> {
Expand Down Expand Up @@ -514,12 +523,14 @@ export class DeeplakeApi {
private async _fetchTables(): Promise<{ tables: string[]; cacheable: boolean }> {
for (let attempt = 0; attempt <= MAX_RETRIES; attempt++) {
try {
const timeoutMs = getQueryTimeoutMs();
const resp = await fetch(`${this.apiUrl}/workspaces/${this.workspaceId}/tables`, {
headers: {
Authorization: `Bearer ${this.token}`,
"X-Activeloop-Org-Id": this.orgId,
...deeplakeClientHeader(),
},
signal: AbortSignal.timeout(timeoutMs),
});
if (resp.ok) {
const data = await resp.json() as { tables?: { table_name: string }[] };
Expand All @@ -533,7 +544,11 @@ export class DeeplakeApi {
continue;
}
return { tables: [], cacheable: false };
} catch {
} catch (e: unknown) {
// A stalled metadata lookup is no more trustworthy than a failed one.
// Return the non-cacheable sentinel immediately, matching query()'s
// per-attempt timeout policy instead of retrying a known deadline.
if (isTimeoutError(e)) return { tables: [], cacheable: false };
if (attempt < MAX_RETRIES) {
await sleep(BASE_DELAY_MS * Math.pow(2, attempt));
continue;
Expand Down Expand Up @@ -755,4 +770,3 @@ export class DeeplakeApi {
export function _resetSdkStateForTesting(): void {
_signalledBalanceExhausted = false;
}

113 changes: 113 additions & 0 deletions tests/shared/deeplake-api-table-timeout.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,113 @@
import { afterEach, describe, expect, it, vi } from "vitest";
import { DeeplakeApi } from "../../src/deeplake-api.js";

const fetchMock = vi.fn();
vi.stubGlobal("fetch", fetchMock);
const nativeAbortTimeout = AbortSignal.timeout.bind(AbortSignal);

function makeApi(): DeeplakeApi {
return new DeeplakeApi("tok", "https://api.test", "org", "ws", "memory");
}

function tablesResponse(...tables: string[]) {
return {
ok: true,
status: 200,
json: async () => ({ tables: tables.map(table_name => ({ table_name })) }),
};
}

async function withWatchdog<T>(operation: Promise<T>): Promise<T> {
let timer: ReturnType<typeof setTimeout> | undefined;
try {
return await Promise.race([
operation,
new Promise<never>((_resolve, reject) => {
timer = setTimeout(() => reject(new Error("test watchdog expired")), 500);
}),
]);
} finally {
if (timer) clearTimeout(timer);
}
}

afterEach(() => {
vi.restoreAllMocks();
fetchMock.mockReset();
delete process.env.HIVEMIND_QUERY_TIMEOUT_MS;
});

describe("DeeplakeApi table discovery timeout", () => {
it.each([
["empty", "", 10_000],
["spaces", " ", 10_000],
["blank whitespace", "\t\r\n", 10_000],
["non-numeric", "abc", 10_000],
["negative", "-1", 10_000],
["fractional", "12.75", 12],
["above the signed timer range", "2147483648", 2_147_483_647],
["infinite", "Infinity", 10_000],
])("normalizes a %s configured timeout and still performs a fast fetch", async (_label, configured, expected) => {
process.env.HIVEMIND_QUERY_TIMEOUT_MS = configured;
const timeoutSpy = vi.spyOn(AbortSignal, "timeout")
.mockImplementation(delay => nativeAbortTimeout(delay));
fetchMock.mockResolvedValueOnce(tablesResponse("memory"));

await expect(makeApi().knownTablesOrNull()).resolves.toEqual(["memory"]);

expect(fetchMock).toHaveBeenCalledOnce();
expect(timeoutSpy).toHaveBeenCalledWith(expected);
});

it.each([
["the documented default", undefined, 10_000],
["zero", "0", 0],
])("preserves %s while a valid fast metadata fetch succeeds", async (_label, configured, expected) => {
if (configured !== undefined) process.env.HIVEMIND_QUERY_TIMEOUT_MS = configured;
const timeoutSpy = vi.spyOn(AbortSignal, "timeout")
.mockImplementation(delay => nativeAbortTimeout(delay));
fetchMock.mockResolvedValueOnce(tablesResponse("memory", "sessions"));

await expect(makeApi().knownTablesOrNull()).resolves.toEqual(["memory", "sessions"]);

expect(timeoutSpy).toHaveBeenCalledWith(expected);
expect(fetchMock.mock.calls[0][1].signal).toBeInstanceOf(AbortSignal);
});

it("uses a native abort signal to stop a metadata fetch that never returns headers", async () => {
process.env.HIVEMIND_QUERY_TIMEOUT_MS = "15";
let abortReason: unknown;
fetchMock.mockImplementation((_url: string, opts: { signal: AbortSignal }) =>
new Promise<never>((_resolve, reject) => {
opts.signal.addEventListener("abort", () => {
abortReason = opts.signal.reason;
reject(opts.signal.reason);
}, { once: true });
}));

await expect(withWatchdog(makeApi().knownTablesOrNull())).resolves.toBeNull();

expect(fetchMock).toHaveBeenCalledOnce();
expect(abortReason).toMatchObject({ name: "TimeoutError" });
});

it("keeps the native signal active while consuming the response body", async () => {
process.env.HIVEMIND_QUERY_TIMEOUT_MS = "15";
let bodyAbortReason: unknown;
fetchMock.mockImplementation((_url: string, opts: { signal: AbortSignal }) => ({
ok: true,
status: 200,
json: () => new Promise<never>((_resolve, reject) => {
opts.signal.addEventListener("abort", () => {
bodyAbortReason = opts.signal.reason;
reject(opts.signal.reason);
}, { once: true });
}),
}));

await expect(withWatchdog(makeApi().knownTablesOrNull())).resolves.toBeNull();

expect(fetchMock).toHaveBeenCalledOnce();
expect(bodyAbortReason).toMatchObject({ name: "TimeoutError" });
});
});