import { hostname } from "node:os"; import { basename } from "node:path"; import { AgentTickClient as AgentTickSdkClient, type CreateRequest, type CreateStatusUpdate, type CreateToolActivity, type PrivateRequestPrepareResponse, type PrivateStatusUpdatePrepareResponse, type RequestRecord, } from "@self-deprecated/agent-tick-sdk"; import { ChoiceFlagSchema, type ChoiceFlag } from "@self-deprecated/agent-tick-shared"; import { redactForAgentTick } from "./redaction"; import type { AgentTickApprovalResponse, AgentTickConfig, AgentTickRequester, CreateApprovalRequestInput, CreateApprovalRequestResult, CreateStatusUpdateInput, CreateToolActivityInput, RequestWaiterStopReason, } from "./types"; import { createPrivateRequestInput } from "./privateStatus"; // Cloudflare's default proxy read timeout is around 100s. Keep each // long-poll comfortably below that and loop client-side for longer waits. const MAX_WAIT_LONG_POLL_MS = 55 * 1000; const INITIAL_WAIT_RETRY_BACKOFF_MS = 1000; const MAX_WAIT_RETRY_BACKOFF_MS = 60 * 1000; const SUPPORTED_CHOICE_FLAGS = new Set(ChoiceFlagSchema.options); class AgentTickHttpError extends Error { constructor(message: string, readonly status: number) { super(message); this.name = "AgentTickHttpError"; } } function isAbortError(error: unknown, signal?: AbortSignal): boolean { if (signal?.aborted) return true; return typeof error === "object" && error !== null && "name" in error && (String((error as { name?: unknown }).name) === "AbortError" || String((error as { name?: unknown }).name) === "TimeoutError"); } function errorStatus(error: unknown): number | undefined { return typeof error === "object" && error !== null && "status" in error && typeof (error as { status?: unknown }).status === "number" ? (error as { status: number }).status : undefined; } function sanitizeError(error: unknown): Error { const message = error instanceof Error ? error.message : String(error); const sanitized = redactForAgentTick(message); const status = errorStatus(error); return status === undefined ? new Error(sanitized) : new AgentTickHttpError(sanitized, status); } function isTransientWaitError(error: unknown, signal?: AbortSignal): boolean { if (isAbortError(error, signal)) return false; const status = errorStatus(error); if (status !== undefined) return status === 408 || status === 429 || status >= 500; // Fetch network failures are usually TypeError in Node/browser runtimes. // The SDK parses JSON before checking non-2xx statuses, so CDN/proxy HTML // timeout bodies can surface as SyntaxError without a status. return error instanceof TypeError || error instanceof SyntaxError; } function waitBackoffDelay(baseMs: number): number { const jitterMultiplier = 0.8 + Math.random() * 0.4; return Math.max(1, Math.round(baseMs * jitterMultiplier)); } function abortReason(signal?: AbortSignal): Error { if (signal?.reason instanceof Error) return signal.reason; return new Error("Agent Tick wait aborted"); } async function sleep(ms: number, signal?: AbortSignal): Promise { if (signal?.aborted) throw abortReason(signal); await new Promise((resolve, reject) => { let timer: ReturnType | undefined; const onAbort = () => { if (timer) clearTimeout(timer); signal?.removeEventListener("abort", onAbort); reject(abortReason(signal)); }; timer = setTimeout(() => { signal?.removeEventListener("abort", onAbort); resolve(); }, ms); signal?.addEventListener("abort", onAbort, { once: true }); if (signal?.aborted) onAbort(); }); } function stringifyMetadataValue(value: unknown): string | undefined { if (value === undefined || value === null) return undefined; if (typeof value === "string") return value; if (typeof value === "number" || typeof value === "boolean" || typeof value === "bigint") return String(value); try { return JSON.stringify(value); } catch { return String(value); } } function normalizeMetadata(metadata: Record | undefined): Record | undefined { if (!metadata) return undefined; const normalized: Record = {}; for (const [key, value] of Object.entries(metadata)) { const stringValue = stringifyMetadataValue(value); if (stringValue !== undefined) normalized[key] = stringValue; } return Object.keys(normalized).length > 0 ? normalized : undefined; } function normalizeRequestType(requestType: CreateApprovalRequestInput["requestType"]): string { if (requestType === "approval") return "sanction"; if (requestType === "steer") return "steering"; return requestType; } function normalizeRequester(input: AgentTickRequester | undefined): CreateRequest["requester"] { return { name: input?.name || "Pi", host: input?.host ?? hostname(), workingDirectory: input?.workingDirectory ?? process.cwd(), clientName: input?.clientName ?? input?.projectName ?? basename(process.cwd()), }; } function isSupportedChoiceFlag(flag: string): flag is ChoiceFlag { return SUPPORTED_CHOICE_FLAGS.has(flag); } function normalizeChoices(input: CreateApprovalRequestInput["choices"]): CreateRequest["choices"] | undefined { if (!input) return undefined; return input.map((choice) => { const flags = choice.flags?.filter(isSupportedChoiceFlag); return { id: choice.id, label: choice.label, ...(choice.description ? { description: choice.description } : {}), ...(choice.kind ? { kind: choice.kind } : {}), ...(flags?.length ? { flags } : {}), ...(choice.tags?.length ? { tags: choice.tags } : {}), }; }); } function normalizeCreateRequest(input: CreateApprovalRequestInput): CreateRequest { return { requester: normalizeRequester(input.requester), requestType: normalizeRequestType(input.requestType), ...(input.sessionId ? { sessionId: input.sessionId } : {}), ...(input.session ? { session: input.session } : {}), title: input.title, ...(input.body ? { body: input.body } : {}), ...(input.command ? { command: input.command } : {}), ...(input.choices ? { choices: normalizeChoices(input.choices) } : {}), ...(input.questions ? { questions: input.questions } : {}), ...(input.allowFreeformReply !== undefined ? { allowFreeformReply: input.allowFreeformReply } : {}), ...(input.metadata ? { metadata: normalizeMetadata(input.metadata) } : {}), }; } function isPrivateRequiredError(error: unknown): boolean { if (errorStatus(error) !== 409) return false; const message = error instanceof Error ? error.message : String(error); return /private requests? (are )?required/i.test(message); } function normalizeCreateResult(value: { request: RequestRecord; waiter?: { token: string; waiterId?: string } }): CreateApprovalRequestResult { const id = value.request.id; if (typeof id !== "string" || !id) throw new Error("Agent Tick response did not include a request id"); return value.waiter?.token ? { id, waiterToken: value.waiter.token, ...(value.waiter.waiterId ? { waiterId: value.waiter.waiterId } : {}) } : { id }; } async function postAgentTickJSON(config: AgentTickConfig, path: string, body: unknown, signal?: AbortSignal): Promise { const response = await fetch(`${config.server.replace(/\/$/, "")}${path}`, { method: "POST", headers: { "Authorization": `Bearer ${config.token}`, "Content-Type": "application/json", }, body: JSON.stringify(body), ...(signal ? { signal } : {}), }); let parsed: any; try { parsed = await response.json(); } catch { parsed = undefined; } if (!response.ok) { const message = typeof parsed?.error?.message === "string" ? parsed.error.message : "Request failed"; throw new AgentTickHttpError(message, response.status); } return parsed as T; } function normalizeStatusInput(input: CreateStatusUpdateInput): CreateStatusUpdate { const metadata = normalizeMetadata(input.metadata); return { message: input.message, state: input.state, ...(input.threadId ? { threadId: input.threadId } : {}), ...(input.sessionId ? { sessionId: input.sessionId } : {}), ...(input.session ? { session: input.session } : {}), ...(input.nextStep ? { nextStep: input.nextStep } : {}), ...(input.host ? { host: input.host } : {}), ...(input.workingDirectory ? { workingDirectory: input.workingDirectory } : {}), ...(input.clientName ?? input.projectName ? { clientName: input.clientName ?? input.projectName } : {}), ...(metadata ? { metadata } : {}), ...(input.contentMode ? { contentMode: input.contentMode } : {}), ...(input.encryptedPayload ? { encryptedPayload: input.encryptedPayload } : {}), ...(input.privateRecipientVersion ? { privateRecipientVersion: input.privateRecipientVersion } : {}), ...(input.contextUsage ? { contextUsage: input.contextUsage } : {}), }; } function normalizeToolActivityInput(input: CreateToolActivityInput): CreateToolActivity { const metadata = normalizeMetadata(input.metadata); return { toolName: input.toolName, state: input.state, ...(input.outcome ? { outcome: input.outcome } : {}), ...(input.threadId ? { threadId: input.threadId } : {}), ...(input.sessionId ? { sessionId: input.sessionId } : {}), ...(input.turnId ? { turnId: input.turnId } : {}), ...(input.toolCallId ? { toolCallId: input.toolCallId } : {}), ...(input.summary ? { summary: input.summary } : {}), ...(metadata ? { metadata } : {}), ...(input.contentMode ? { contentMode: input.contentMode } : {}), ...(input.encryptedPayload ? { encryptedPayload: input.encryptedPayload } : {}), ...(input.privateRecipientVersion ? { privateRecipientVersion: input.privateRecipientVersion } : {}), ...(input.startedAt ? { startedAt: input.startedAt } : {}), ...(input.finishedAt ? { finishedAt: input.finishedAt } : {}), }; } function normalizeWaitResult(value: { request: RequestRecord; terminal: boolean }): AgentTickApprovalResponse { return { requestId: value.request.id, status: value.request.status, state: value.request.status, terminal: value.terminal, ...value.request.response, response: value.request.response, }; } function isPendingWait(value: { request: RequestRecord; terminal: boolean }): boolean { return value.terminal === false || value.request.status === "pending"; } function fetchWithSignal(signal: AbortSignal): typeof fetch { return ((input: RequestInfo | URL, init?: RequestInit) => fetch(input, { ...init, signal: init?.signal ?? signal, })) as typeof fetch; } export class AgentTickClient { constructor(private readonly config: AgentTickConfig) {} private sdk(signal?: AbortSignal): AgentTickSdkClient { return new AgentTickSdkClient({ baseUrl: this.config.server, tokenProvider: () => this.config.token, ...(signal ? { fetch: fetchWithSignal(signal) } : {}), }); } async createApprovalRequest(input: CreateApprovalRequestInput, signal?: AbortSignal): Promise { const requestInput = normalizeCreateRequest(input); try { return normalizeCreateResult(await postAgentTickJSON<{ request: RequestRecord; waiter?: { token: string; waiterId?: string } }>(this.config, "/v1/requests", requestInput, signal)); } catch (error) { if (!isPrivateRequiredError(error)) throw sanitizeError(error); try { const prepare = await this.preparePrivateRequest({ requestType: requestInput.requestType }, signal); return normalizeCreateResult(await postAgentTickJSON<{ request: RequestRecord; waiter?: { token: string; waiterId?: string } }>(this.config, "/v1/requests", createPrivateRequestInput(requestInput, prepare), signal)); } catch (privateError) { throw sanitizeError(privateError); } } } async preparePrivateRequest(input: { requestType?: CreateRequest["requestType"] } = {}, signal?: AbortSignal): Promise { try { return await this.sdk(signal).preparePrivateRequest(input); } catch (error) { throw sanitizeError(error); } } async waitForApproval(id: string, waiterToken?: string, timeoutMs = 0, signal?: AbortSignal): Promise { const unlimited = timeoutMs <= 0; const deadline = unlimited ? Number.POSITIVE_INFINITY : Date.now() + timeoutMs; const waitSdk = this.sdk(signal); let retryBackoffMs = INITIAL_WAIT_RETRY_BACKOFF_MS; while (true) { try { const remainingMs = unlimited ? MAX_WAIT_LONG_POLL_MS : Math.max(1, deadline - Date.now()); const effectiveTimeoutMs = Math.min(MAX_WAIT_LONG_POLL_MS, remainingMs); const parsed = await waitSdk.waitForRequest(id, { timeoutMs: effectiveTimeoutMs, ...(waiterToken ? { waiterToken } : {}), ...(signal ? { signal } : {}), }); retryBackoffMs = INITIAL_WAIT_RETRY_BACKOFF_MS; if (!isPendingWait(parsed)) return normalizeWaitResult(parsed); if (!unlimited && Date.now() >= deadline) return normalizeWaitResult(parsed); } catch (error) { if (!isTransientWaitError(error, signal) || (!unlimited && Date.now() >= deadline)) throw sanitizeError(error); const remainingMs = unlimited ? Number.POSITIVE_INFINITY : Math.max(0, deadline - Date.now()); const delayMs = Math.min(waitBackoffDelay(retryBackoffMs), remainingMs); if (delayMs <= 0) throw sanitizeError(error); await sleep(delayMs, signal); retryBackoffMs = Math.min(MAX_WAIT_RETRY_BACKOFF_MS, retryBackoffMs * 2); } } } async abandonApproval(id: string, signal?: AbortSignal): Promise { try { await this.sdk(signal).abandonRequest(id); } catch (error) { throw sanitizeError(error); } } async stopRequestWaiter(id: string, waiterToken: string, reason: RequestWaiterStopReason, signal?: AbortSignal): Promise { try { await this.sdk(signal).stopRequestWaiter(id, { reason }, { waiterToken }); } catch (error) { throw sanitizeError(error); } } async reportRequestWaiterError(id: string, waiterToken: string, code: string, message?: string, signal?: AbortSignal): Promise { try { await this.sdk(signal).reportRequestWaiterError(id, { code, ...(message ? { message } : {}) }, { waiterToken }); } catch (error) { throw sanitizeError(error); } } async preparePrivateStatusUpdate(signal?: AbortSignal): Promise { try { return await this.sdk(signal).preparePrivateStatusUpdate(); } catch (error) { throw sanitizeError(error); } } async createStatusUpdate(input: CreateStatusUpdateInput, signal?: AbortSignal): Promise { try { await postAgentTickJSON(this.config, "/v1/status-updates", normalizeStatusInput(input), signal); } catch (error) { throw sanitizeError(error); } } async createToolActivity(input: CreateToolActivityInput, signal?: AbortSignal): Promise { try { await postAgentTickJSON(this.config, "/v1/tool-activities", normalizeToolActivityInput(input), signal); } catch (error) { throw sanitizeError(error); } } }