> ## Documentation Index
> Fetch the complete documentation index at: https://docs.oao.sh/llms.txt
> Use this file to discover all available pages before exploring further.

# Connect a codebase

> Give one OAO agent safe, durable access to repository context from a server-side integration.

Connect a codebase by keeping the OAO credential and repository access in a
trusted server process. Create the first session and run together, execute any
caller-owned file tools under a narrow repository policy, and use the durable
project event stream to resume after disconnects.

The current MVP supports three code-context patterns:

1. Attach up to eight supported files directly to an initial or later turn.
   OAO copies their original bytes into the agent sandbox. See
   [Send files to an agent](/integrations/files) for media types and limits.
2. Put a small, trusted text excerpt directly in `initialMessage` or a later
   turn. The complete message must be no more than 100,000 JavaScript
   characters.
3. Publish narrow `caller` tools such as `read_repository_file` and
   `search_repository`, then run those tools in your own server process against
   one allowlisted repository root.

Use attached files for a small, deliberate set whose filenames should remain
visible in the transcript. Use inline text for a tiny excerpt. Use caller tools
when the agent needs to choose which files to inspect over several turns or when
the repository exceeds the per-turn attachment limits. Caller-tool arguments
and results are durable public records, so return only content that is safe to
persist and display.

## Create a server-side API key

Open **API keys** in the console, select **Create API key**, enter a descriptive
name, and grant only the scopes the integration needs. Save the secret from the
one-time acknowledgement dialog before closing it.

For automated provisioning, create the integration key through the API using an
already-authenticated human or service principal with `project:admin`:

```bash theme={null}
curl --fail-with-body \
  --request POST \
  --header "Authorization: Bearer $OAO_OPERATOR_TOKEN" \
  --header 'Content-Type: application/json' \
  --header 'Idempotency-Key: create-codebase-integration-key-1' \
  --data '{
    "name": "Codebase integration",
    "scopes": [
      "agent:write",
      "session:read",
      "session:write",
      "run:create",
      "run:read",
      "tool_call:claim",
      "tool_call:submit"
    ]
  }' \
  "$OAO_API_URL/v1/projects/$OAO_PROJECT_ID/api-keys"
```

The successful creation response shows the `oao_...` secret once. Copy it into
the integration's secret store immediately. OAO stores only a keyed hash and
cannot reveal the plaintext again. See [Authentication](/reference/authentication)
for browser-cookie authentication, tenant authorization, and key handling.

Grant the smallest set of scopes for the process you are running:

| Capability | Required scopes |
| - | - |
| Create an initial session and watch it | `session:read`, `session:write`, `run:create`, `run:read` |
| Send later turns | `run:create` |
| Execute caller-owned tools | `tool_call:claim`, `tool_call:submit` |
| Create the agent in the SDK example below | `agent:write` |
| Request cancellation during shutdown | `run:cancel` |

The complete SDK example uses every scope except `run:cancel`. An operator can
create the agent separately and remove `agent:write` from the long-lived worker
key afterward.

Store the key only in the server environment:

```bash theme={null}
export OAO_API_URL=http://127.0.0.1:3000
export OAO_API_KEY='oao_...'
export OAO_PROJECT_ID='00000000-0000-4000-8000-000000000002'
```

`OAO_API_URL` is the API origin, without `/v1`, when using `@oao/sdk-js`.
External HTTP examples below add `/v1` themselves.

<Warning>
  Never put an OAO API key in browser JavaScript, a mobile bundle, a repository,
  a query string, an SSE URL, or logs. Send it only as `Authorization: Bearer
      oao_...` from the server.
</Warning>

Every mutation also needs an `Idempotency-Key`. Persist a key before the first
attempt, and reuse that same key and identical JSON body when the outcome is
uncertain. A replay returns the original public response. Reusing the key for a
different body returns `409 idempotency_conflict`.

## Pattern 1: attach a small trusted text file

Read the selected file on your server and send canonical base64. The API stores
the original bytes durably with the run, verifies them by SHA-256 at dispatch,
copies them into the sandbox, and exposes only safe metadata in transcript
reads. It does not inject the file content into the model prompt:

```ts theme={null}
const source = await readFile(operatorSelectedPath);
const created = await client.createSession(
  projectId,
  {
    agentId,
    title: "Repository review",
    initialMessage: "Review the attached repository entry point.",
    files: [
      {
        name: repositoryRelativePath.split("/").at(-1)!,
        contentType: "application/typescript",
        dataBase64: source.toString("base64"),
      },
    ],
  },
  { idempotencyKey: persistedCreateSessionKey },
);
```

Do not use this pattern for secrets, entire repositories, or unsupported
binaries. File names cannot contain paths, so use a plain display filename and
put any reviewed repository-relative path in the message.

## Pattern 2: inline a small trusted text excerpt

Read and review the text on your server before including it. Reject binary
data and secrets, add a clear path label, and check the final assembled message
rather than only the file length:

````ts theme={null}
const source = await readFile(operatorSelectedPath, "utf8");
const initialMessage = [
  "Review this trusted repository excerpt. Focus on public API boundaries.",
  `Path: ${repositoryRelativePath}`,
  "```text",
  source,
  "```",
].join("\n");

if (initialMessage.length > 100_000) {
  throw new Error("The complete initialMessage exceeds 100,000 characters");
}

const created = await client.createSession(
  projectId,
  { agentId, title: "Repository review", initialMessage },
  { idempotencyKey: persistedCreateSessionKey },
);
````

This is text in the conversation, not an attachment. Do not use this pattern
for untrusted files, generated archives, images, PDFs, executables, entire
repositories, or any content that may contain credentials.

## Pattern 3: run allowlisted repository tools

The following private-workspace example creates one agent, creates its first
session atomically with `initialMessage`, handles resumable SSE, services two
caller-owned tools, and optionally submits one later turn after the previous
run settles.

Add the private SDK to an OAO workspace package that owns the integration:

```bash theme={null}
pnpm --filter <your-server-package> add @oao/sdk-js@workspace:*
```

The SDK is currently `private: true`; this install form is for packages in this
monorepo. Save the following as `src/codebase-agent.ts` in that package and run
it with the package's TypeScript runner.

```ts theme={null}
import { randomUUID } from "node:crypto";
import {
  lstat,
  open,
  readdir,
  readFile,
  realpath,
  rename,
  writeFile,
} from "node:fs/promises";
import { extname, isAbsolute, relative, resolve, sep } from "node:path";
import {
  OaoApiError,
  OaoClient,
  type ProductEvent,
  type Run,
  type ToolCall,
} from "@oao/sdk-js";

const required = (name: string): string => {
  const value = process.env[name];
  if (!value) throw new Error(`${name} is required`);
  return value;
};

const apiUrl = required("OAO_API_URL");
const apiKey = required("OAO_API_KEY");
const projectId = required("OAO_PROJECT_ID");
const repositoryRoot = await realpath(required("OAO_REPOSITORY_ROOT"));
const statePath = resolve(
  process.env.OAO_STATE_FILE ?? ".oao-codebase-agent-state.json",
);
const followUp = process.env.OAO_FOLLOW_UP;

const client = new OaoClient({ baseUrl: apiUrl, apiKey });
const shutdown = new AbortController();
const stop = () => shutdown.abort();
process.once("SIGINT", stop);
process.once("SIGTERM", stop);

const TERMINAL_STATES = new Set<Run["state"]>([
  "completed",
  "failed",
  "cancelled",
  "timed_out",
]);
const ALLOWED_EXTENSIONS = new Set([
  ".c",
  ".cc",
  ".css",
  ".go",
  ".h",
  ".html",
  ".java",
  ".js",
  ".json",
  ".jsx",
  ".md",
  ".mdx",
  ".py",
  ".rb",
  ".rs",
  ".sql",
  ".swift",
  ".toml",
  ".ts",
  ".tsx",
  ".txt",
  ".yaml",
  ".yml",
]);
const DENIED_SEGMENTS = new Set([
  ".git",
  ".hg",
  ".svn",
  ".ssh",
  "dist",
  "node_modules",
  "vendor",
]);
const DENIED_FILE = /^(?:\.env(?:\..*)?|credentials?|id_[a-z0-9_-]+)$/iu;
const DENIED_EXTENSION = /\.(?:key|p12|pfx|pem)$/iu;
const SECRET_CONTENT =
  /-----BEGIN [A-Z ]*PRIVATE KEY-----|(?:api[_-]?key|client[_-]?secret|password|private[_-]?key|refresh[_-]?token)\s*[:=]\s*["']?[A-Za-z0-9_./+=-]{12,}/iu;
const MAX_FILE_BYTES = 256 * 1024;
const MAX_SEARCHED_FILES = 2_000;
const MAX_MATCHES = 50;

interface State {
  agentId?: string;
  agentCreateKey?: string;
  sessionId?: string;
  initialRunId?: string;
  sessionCreateKey?: string;
  followUpRunId?: string;
  followUpKey?: string;
  lastEventId?: string;
}

async function loadState(): Promise<State> {
  try {
    return JSON.parse(await readFile(statePath, "utf8")) as State;
  } catch (error) {
    if ((error as NodeJS.ErrnoException).code === "ENOENT") return {};
    throw error;
  }
}

let state = await loadState();

async function saveState(): Promise<void> {
  const temporary = `${statePath}.${process.pid}.tmp`;
  await writeFile(temporary, `${JSON.stringify(state, null, 2)}\n`, {
    mode: 0o600,
  });
  await rename(temporary, statePath);
}

function isContained(path: string): boolean {
  const fromRoot = relative(repositoryRoot, path);
  return (
    fromRoot === "" ||
    (!fromRoot.startsWith(`..${sep}`) &&
      fromRoot !== ".." &&
      !isAbsolute(fromRoot))
  );
}

function validateRelativePath(input: string): string {
  if (!input || input.includes("\0") || isAbsolute(input)) {
    throw new InputError();
  }
  const candidate = resolve(repositoryRoot, input);
  if (!isContained(candidate)) throw new InputError();
  const segments = relative(repositoryRoot, candidate).split(sep);
  if (
    segments.some((segment) =>
      DENIED_SEGMENTS.has(segment.toLocaleLowerCase("en-US")),
    ) ||
    DENIED_FILE.test(segments.at(-1) ?? "") ||
    DENIED_EXTENSION.test(segments.at(-1) ?? "")
  ) {
    throw new InputError();
  }
  return candidate;
}

async function assertNoSymlinkComponents(path: string): Promise<void> {
  const parts = relative(repositoryRoot, path).split(sep).filter(Boolean);
  let current = repositoryRoot;
  for (const part of parts) {
    current = resolve(current, part);
    if ((await lstat(current)).isSymbolicLink()) throw new InputError();
  }
}

async function readSafeText(relativePath: string): Promise<string> {
  const candidate = validateRelativePath(relativePath);
  await assertNoSymlinkComponents(candidate);
  const canonical = await realpath(candidate);
  if (!isContained(canonical)) throw new InputError();
  if (!ALLOWED_EXTENSIONS.has(extname(canonical).toLocaleLowerCase("en-US"))) {
    throw new InputError();
  }

  const handle = await open(canonical, "r");
  try {
    const metadata = await handle.stat();
    if (!metadata.isFile() || metadata.size > MAX_FILE_BYTES) {
      throw new InputError();
    }
    const bytes = await handle.readFile();
    if (bytes.length > MAX_FILE_BYTES || bytes.includes(0)) {
      throw new InputError();
    }
    let text: string;
    try {
      text = new TextDecoder("utf-8", { fatal: true }).decode(bytes);
    } catch {
      throw new InputError();
    }
    if (SECRET_CONTENT.test(text)) throw new InputError();
    return text;
  } finally {
    await handle.close();
  }
}

function publicPath(absolutePath: string): string {
  return relative(repositoryRoot, absolutePath).split(sep).join("/");
}

async function searchRepository(query: string) {
  if (!query || query.length > 200 || /[\u0000-\u001f]/u.test(query)) {
    throw new InputError();
  }
  const needle = query.toLocaleLowerCase("en-US");
  const matches: Array<{ path: string; line: number; excerpt: string }> = [];
  let inspected = 0;

  const visit = async (directory: string): Promise<void> => {
    const entries = (await readdir(directory, { withFileTypes: true })).sort(
      (left, right) => left.name.localeCompare(right.name),
    );
    for (const entry of entries) {
      if (matches.length >= MAX_MATCHES || inspected >= MAX_SEARCHED_FILES) {
        return;
      }
      if (entry.isSymbolicLink()) continue;
      const name = entry.name.toLocaleLowerCase("en-US");
      if (entry.isDirectory()) {
        if (!DENIED_SEGMENTS.has(name))
          await visit(resolve(directory, entry.name));
        continue;
      }
      if (!entry.isFile()) continue;
      inspected += 1;
      const absolutePath = resolve(directory, entry.name);
      const path = publicPath(absolutePath);
      try {
        const text = await readSafeText(path);
        const lines = text.split(/\r?\n/u);
        for (let index = 0; index < lines.length; index += 1) {
          if (lines[index]!.toLocaleLowerCase("en-US").includes(needle)) {
            matches.push({
              path,
              line: index + 1,
              excerpt: lines[index]!.slice(0, 500),
            });
            if (matches.length >= MAX_MATCHES) break;
          }
        }
      } catch (error) {
        if (!(error instanceof InputError)) throw error;
      }
    }
  };

  await visit(repositoryRoot);
  return { matches };
}

class InputError extends Error {}

function objectArgument(value: unknown): Record<string, unknown> {
  if (!value || Array.isArray(value) || typeof value !== "object") {
    throw new InputError();
  }
  return value as Record<string, unknown>;
}

async function executeTool(call: ToolCall) {
  try {
    const input = objectArgument(call.safeArguments);
    if (call.toolName === "read_repository_file") {
      if (typeof input.path !== "string") throw new InputError();
      return {
        version: 1,
        status: "success",
        value: {
          path: publicPath(validateRelativePath(input.path)),
          content: await readSafeText(input.path),
        },
      } as const;
    }
    if (call.toolName === "search_repository") {
      if (typeof input.query !== "string") throw new InputError();
      return {
        version: 1,
        status: "success",
        value: await searchRepository(input.query),
      } as const;
    }
    throw new InputError();
  } catch (error) {
    return {
      version: 1,
      status: "failure",
      error: {
        code:
          error instanceof InputError
            ? "invalid_tool_arguments"
            : "tool_failed",
        message:
          error instanceof InputError
            ? "The requested repository input is not available"
            : "Repository tool execution failed",
      },
    } as const;
  }
}

const sleep = (milliseconds: number, signal: AbortSignal) =>
  new Promise<void>((resolveSleep, reject) => {
    const timer = setTimeout(resolveSleep, milliseconds);
    signal.addEventListener(
      "abort",
      () => {
        clearTimeout(timer);
        reject(signal.reason ?? new DOMException("Aborted", "AbortError"));
      },
      { once: true },
    );
  });

function isAbort(error: unknown): boolean {
  return error instanceof DOMException && error.name === "AbortError";
}

async function retryMutation<T>(operation: () => Promise<T>): Promise<T> {
  let lastError: unknown;
  for (const delay of [0, 500, 1_500]) {
    if (delay) await sleep(delay, shutdown.signal);
    try {
      return await operation();
    } catch (error) {
      lastError = error;
      const retryable =
        error instanceof TypeError ||
        (error instanceof OaoApiError &&
          (error.status === 429 || error.status >= 500));
      if (!retryable) throw error;
    }
  }
  throw lastError;
}

async function serviceToolCall(call: ToolCall): Promise<void> {
  const claimKey = randomUUID();
  let claim;
  try {
    claim = await retryMutation(() =>
      client.claimToolCall(
        projectId,
        call.id,
        { leaseMs: 60_000 },
        { idempotencyKey: claimKey, signal: shutdown.signal },
      ),
    );
  } catch (error) {
    if (!isAbort(error)) console.warn(`Tool ${call.id} was not claimable`);
    return;
  }

  const heartbeatAbort = new AbortController();
  const heartbeatSignal = AbortSignal.any([
    shutdown.signal,
    heartbeatAbort.signal,
  ]);
  let leaseError: unknown;
  const heartbeat = (async () => {
    try {
      while (true) {
        await sleep(20_000, heartbeatSignal);
        const renewalKey = randomUUID();
        await retryMutation(() =>
          client.renewToolCall(
            projectId,
            call.id,
            { fence: claim.fence, leaseMs: 60_000 },
            { idempotencyKey: renewalKey, signal: heartbeatSignal },
          ),
        );
      }
    } catch (error) {
      if (!isAbort(error)) leaseError = error;
    }
  })();

  const safeResult = await executeTool(call);
  heartbeatAbort.abort();
  await heartbeat;
  if (leaseError || shutdown.signal.aborted) return;

  const resultKey = randomUUID();
  try {
    await retryMutation(() =>
      client.submitToolResult(
        projectId,
        call.id,
        { fence: claim.fence, safeResult },
        { idempotencyKey: resultKey, signal: shutdown.signal },
      ),
    );
  } catch (error) {
    // Do not submit with a new fence or key. A stale/expired claim must be
    // reclaimed, while an uncertain retry must retain resultKey and the body.
    if (!isAbort(error))
      console.error(`Tool ${call.id} result was not accepted`);
  }
}

const inFlight = new Set<string>();

async function drainCallerTools(runId: string): Promise<void> {
  const page = await client.listToolCalls(
    projectId,
    { runId, limit: 200 },
    { signal: shutdown.signal },
  );
  for (const call of page.data) {
    if (
      call.owner !== "caller" ||
      !["caller_pending", "caller_claimed"].includes(call.stage) ||
      inFlight.has(call.id)
    ) {
      continue;
    }
    inFlight.add(call.id);
    try {
      await serviceToolCall(call);
    } finally {
      inFlight.delete(call.id);
    }
  }
}

async function applyEvent(event: ProductEvent, runId: string): Promise<void> {
  if (
    event.kind === "tool_call.requested" ||
    event.kind === "tool_call.claimed" ||
    event.kind === "run.state_changed"
  ) {
    await drainCallerTools(runId);
  }
}

async function waitForSettlement(runId: string): Promise<Run> {
  await drainCallerTools(runId);
  let run = await client.getRun(projectId, runId, {
    signal: shutdown.signal,
  });
  if (TERMINAL_STATES.has(run.state)) return run;

  const streamAbort = new AbortController();
  const streamSignal = AbortSignal.any([shutdown.signal, streamAbort.signal]);
  try {
    for await (const frame of client.streamProjectEvents(projectId, {
      lastEventId: state.lastEventId,
      reconnect: true,
      signal: streamSignal,
    })) {
      await applyEvent(frame.data, runId);
      if (frame.id) {
        state = { ...state, lastEventId: frame.id };
        await saveState();
      }
      run = await client.getRun(projectId, runId, { signal: streamSignal });
      if (TERMINAL_STATES.has(run.state)) return run;
    }
    throw new Error("The project event stream ended unexpectedly");
  } finally {
    streamAbort.abort();
  }
}

const toolConfiguration = [
  {
    name: "read_repository_file",
    description:
      "Read one small UTF-8 text file from the allowlisted repository.",
    owner: "caller" as const,
    approval: "never" as const,
    inputSchema: {
      type: "object",
      properties: { path: { type: "string" } },
      required: ["path"],
      additionalProperties: false,
    },
    outputSchema: {
      type: "object",
      properties: {
        path: { type: "string" },
        content: { type: "string" },
      },
      required: ["path", "content"],
      additionalProperties: false,
    },
  },
  {
    name: "search_repository",
    description:
      "Search small allowlisted UTF-8 repository files for literal text.",
    owner: "caller" as const,
    approval: "never" as const,
    inputSchema: {
      type: "object",
      properties: { query: { type: "string" } },
      required: ["query"],
      additionalProperties: false,
    },
    outputSchema: {
      type: "object",
      properties: {
        matches: {
          type: "array",
          items: {
            type: "object",
            properties: {
              path: { type: "string" },
              line: { type: "integer" },
              excerpt: { type: "string" },
            },
            required: ["path", "line", "excerpt"],
            additionalProperties: false,
          },
        },
      },
      required: ["matches"],
      additionalProperties: false,
    },
  },
] as const;

async function main(): Promise<void> {
  if (!state.agentId) {
    state.agentCreateKey ??= randomUUID();
    await saveState();
    const agent = await retryMutation(() =>
      client.createAgent(
        projectId,
        {
          key: "repository-reviewer",
          name: "Repository reviewer",
          description:
            "Reviews one allowlisted repository through caller tools.",
          initialConfig: {
            systemPrompt:
              "Review code using only the provided repository tools. Never request or expose credentials, secrets, raw reasoning, or files outside the allowlisted repository.",
            modelPreset: "project-model-v1",
            tools: toolConfiguration,
            sandbox: {
              enabled: false,
              provider: "daytona-primary",
              network: "none",
              capabilities: ["filesystem_read", "filesystem_write", "shell"],
            },
            limits: { maxTurns: 32, timeoutMs: 60_000 },
          },
        },
        {
          idempotencyKey: state.agentCreateKey!,
          signal: shutdown.signal,
        },
      ),
    );
    state = { ...state, agentId: agent.id };
    await saveState();
  }

  if (!state.sessionId || !state.initialRunId) {
    state.sessionCreateKey ??= randomUUID();
    await saveState();
    const created = await retryMutation(() =>
      client.createSession(
        projectId,
        {
          agentId: state.agentId!,
          title: "Repository review",
          initialMessage:
            process.env.OAO_INITIAL_MESSAGE ??
            "Inspect the repository entry points and summarize the public architecture.",
        },
        {
          idempotencyKey: state.sessionCreateKey!,
          signal: shutdown.signal,
        },
      ),
    );
    state = {
      ...state,
      sessionId: created.id,
      initialRunId: created.latestRunId,
    };
    await saveState();
  }

  let settled = await waitForSettlement(state.initialRunId!);
  console.log(`Initial run ${settled.id}: ${settled.state}`);

  if (followUp && !state.followUpRunId) {
    // waitForSettlement proved the session's latest run is terminal before this
    // mutation. OAO also enforces that invariant transactionally.
    state.followUpKey ??= randomUUID();
    await saveState();
    const next = await retryMutation(() =>
      client.submitRun(
        projectId,
        state.sessionId!,
        { redactedInput: followUp },
        {
          idempotencyKey: state.followUpKey!,
          signal: shutdown.signal,
        },
      ),
    );
    state = { ...state, followUpRunId: next.id };
    await saveState();
  }

  if (state.followUpRunId) {
    settled = await waitForSettlement(state.followUpRunId);
    console.log(`Follow-up run ${settled.id}: ${settled.state}`);
  }

  const session = await client.getSession(projectId, state.sessionId!, {
    signal: shutdown.signal,
  });
  console.log(JSON.stringify(session, null, 2));
}

try {
  await main();
} catch (error) {
  if (!shutdown.signal.aborted) throw error;
  console.log("Shutdown requested; the saved event cursor will be reused.");
} finally {
  shutdown.abort();
  process.removeListener("SIGINT", stop);
  process.removeListener("SIGTERM", stop);
}
```

Run it with a secret-free checkout dedicated to this worker:

```bash theme={null}
export OAO_REPOSITORY_ROOT=/absolute/path/to/allowlisted/repository
export OAO_INITIAL_MESSAGE='Find the HTTP entry point and explain its boundaries.'
export OAO_FOLLOW_UP='Now identify the tests that cover those boundaries.'
pnpm --filter <your-server-package> exec tsx src/codebase-agent.ts
```

The state file contains identifiers, idempotency keys, and the last applied
event cursor, but not the API secret. Put it outside the reviewed repository in
production, restrict its file permissions, and back it with durable storage if
the worker can move between hosts.

<Warning>
  This runnable example persists session and run mutations, but keeps each tool
  claim, renewal, and result idempotency key only for the life of the process.
  Before using the worker for effectful production tools, add the per-tool
  operation journal described below. Repository reads are safe to repeat, but
  arbitrary external side effects may not be.
</Warning>

### Why the file boundary is strict

The example applies all of these checks before returning content:

* Resolves one configured repository root to a canonical absolute path.
* Rejects absolute paths, NUL bytes, traversal outside the root, denied
  directories, secret-like filenames, key/certificate extensions, and unknown
  text extensions.
* Rejects every symbolic-link component and verifies the canonical result is
  still under the root.
* Opens only regular files, caps each read at 256 KiB, rejects NUL-containing or
  invalid UTF-8 data, and applies a conservative secret-content check.
* Treats search as literal text, caps files inspected, matches returned, and
  excerpt length, and skips unreadable or disallowed files.
* Returns generic failure text instead of exception messages, absolute paths,
  file metadata, or stack traces.

These are defense-in-depth controls, not a universal secret detector. Use a
read-only, secret-free checkout; run the worker under a dedicated OS identity;
keep it separate from credential directories; and do not allow an untrusted
process to mutate the checkout while it is being read. Tighten the extension,
directory, size, and content policies for your repository.

## Claim leases, fences, and results

A caller tool starts at `caller_pending`. The worker claims it for a lease from
1,000 through 300,000 milliseconds and receives a positive-integer string
`fence`. Renew and release operations must use the current fence. The example
renews a 60-second lease every 20 seconds and stops the heartbeat before result
submission.

Submit the fence with exactly one safe, versioned result envelope:

```json theme={null}
{
  "fence": "4",
  "safeResult": {
    "version": 1,
    "status": "success",
    "value": {
      "path": "src/server.ts",
      "content": "..."
    }
  }
}
```

The `value` must match the published `outputSchema`. A safe failure uses
`status: "failure"` with one of the current public failure codes and a generic
message. Never put credentials, raw tool payloads, absolute host paths, stack
traces, or raw model reasoning in `safeResult`.

Use one idempotency key per logical claim, renewal, release, or result. Retry an
uncertain request with the same key and body. A tool result is immutable: the
same result key and body returns `outcome: "replayed"`; a different key or body
conflicts. If a lease expires or its fence becomes stale, do not submit anyway.
Let the work be reclaimed and re-executed under the new fence.

### Make production tool processing crash-safe

Persist one operation record per `toolCallId` in your own durable store. At a
minimum it needs `phase`, `claimIdempotencyKey`, the claimed `fence`, the exact
result envelope (or its canonical bytes and hash), and
`resultIdempotencyKey`. Update this record before each network request:

1. Generate and persist the claim key and claim body, then call `claim`.
2. Persist the returned fence before executing the handler.
3. For an effectful handler, pass a stable downstream idempotency key derived
   from the tool-call identity; a lease fence prevents stale OAO commits but
   cannot undo an external effect that already started.
4. Persist the final safe-result body and result key before submitting it.
5. Mark the operation complete only after OAO accepts or replays that exact
   result.

After a crash, retry an uncertain request with its stored key and identical
body. If the stored fence has expired, re-list the tool call and reclaim it;
never submit a result under the stale fence. This journal is what extends OAO's
durable tool ledger across crashes in your integration process.

## Resume the event stream durably

`GET /v1/projects/{projectId}/events` is project-wide. Each SSE frame's `id` is
an opaque cursor backed by committed PostgreSQL position. Persist it only after
your application has applied the event, then reconnect with:

```http theme={null}
Last-Event-ID: djE6MTQy
```

The SDK does this automatically when `reconnect: true`. It waits
`reconnectDelayMs` (1 second by default) between connections, or an SSE `retry`
value if the server sends one. The server closes a normal stream about every 25
seconds; that is not data loss. It also writes `:` comment lines when the
stream opens and after 10 seconds of silence, so proxies keep the connection
alive. Parsers must skip those lines. PostgreSQL is authoritative, while
`LISTEN/NOTIFY` only wakes readers.

Expect duplicate events after a crash between applying an event and saving its
cursor. Make your projection idempotent by event `id` or by
`(aggregateType, aggregateId, aggregateSequence)`. Re-fetch the run or session
when an event kind is unknown. The example always re-reads the target run and
lists durable tool calls, so an SSE frame is a notification rather than the
sole copy of state.

On `SIGINT` or `SIGTERM`, abort the stream and any outstanding HTTP request,
stop lease renewal, finish the state-file rename, and exit. On restart, load
the saved cursor and any persisted mutation key before making another request.
If a process dies while holding a tool claim, its lease expires and another
worker can claim the durable work.

## External repository: use HTTP directly

Because `@oao/sdk-js` is not currently published, a server outside the OAO
monorepo should use the HTTP and SSE contracts. This complete Node.js 22
example uses the inline-text pattern, creates a session, watches its first run,
and sends one optional follow-up only after settlement. Use it with an agent
that has no caller-owned tools.

````ts theme={null}
import { readFile, rename, writeFile } from "node:fs/promises";
import { resolve } from "node:path";

const required = (name: string): string => {
  const value = process.env[name];
  if (!value) throw new Error(`${name} is required`);
  return value;
};

const origin = required("OAO_API_URL").replace(/\/+$/u, "");
const apiKey = required("OAO_API_KEY");
const projectId = required("OAO_PROJECT_ID");
const agentId = required("OAO_AGENT_ID");
const contextPath = resolve(required("OAO_TRUSTED_TEXT_FILE"));
const statePath = resolve(process.env.OAO_STATE_FILE ?? ".oao-http-state.json");
const shutdown = new AbortController();
const stop = () => shutdown.abort();
process.once("SIGINT", stop);
process.once("SIGTERM", stop);

interface State {
  createSessionKey?: string;
  sessionId?: string;
  runId?: string;
  followUpKey?: string;
  followUpRunId?: string;
  lastEventId?: string;
}

async function loadState(): Promise<State> {
  try {
    return JSON.parse(await readFile(statePath, "utf8")) as State;
  } catch (error) {
    if ((error as NodeJS.ErrnoException).code === "ENOENT") return {};
    throw error;
  }
}

let state = await loadState();

async function saveState(): Promise<void> {
  const temporary = `${statePath}.${process.pid}.tmp`;
  await writeFile(temporary, `${JSON.stringify(state, null, 2)}\n`, {
    mode: 0o600,
  });
  await rename(temporary, statePath);
}

async function request<T>(
  path: string,
  init: RequestInit = {},
  idempotencyKey?: string,
): Promise<T> {
  const headers = new Headers(init.headers);
  headers.set("authorization", `Bearer ${apiKey}`);
  headers.set("accept", "application/json");
  if (init.body) headers.set("content-type", "application/json");
  if (idempotencyKey) headers.set("idempotency-key", idempotencyKey);
  const response = await fetch(`${origin}/v1${path}`, {
    ...init,
    headers,
    signal: shutdown.signal,
  });
  if (!response.ok) {
    const body = (await response.json().catch(() => undefined)) as
      | { error?: { code?: string; message?: string; requestId?: string } }
      | undefined;
    throw new Error(
      `${body?.error?.code ?? response.status}: ${body?.error?.message ?? "OAO request failed"} (${body?.error?.requestId ?? response.headers.get("x-request-id") ?? "no request id"})`,
    );
  }
  return (await response.json()) as T;
}

interface Run {
  id: string;
  state:
    | "queued"
    | "running"
    | "waiting_for_tool"
    | "waiting_for_approval"
    | "retry_scheduled"
    | "completed"
    | "failed"
    | "cancelled"
    | "timed_out";
}

interface CreatedSession {
  id: string;
  latestRunId: string;
  run: Run;
}

interface SseFrame {
  id?: string;
  event?: string;
  data: string;
  retry?: number;
}

async function* parseSse(
  body: ReadableStream<Uint8Array>,
): AsyncGenerator<SseFrame> {
  const reader = body.getReader();
  const decoder = new TextDecoder();
  let buffer = "";
  let data: string[] = [];
  let id: string | undefined;
  let event: string | undefined;
  let retry: number | undefined;
  const consume = (line: string): SseFrame | undefined => {
    if (line === "") {
      if (!data.length) return undefined;
      const frame = { id, event, retry, data: data.join("\n") };
      data = [];
      id = undefined;
      event = undefined;
      retry = undefined;
      return frame;
    }
    if (line.startsWith(":")) return undefined;
    const separator = line.indexOf(":");
    const field = separator < 0 ? line : line.slice(0, separator);
    const value = (separator < 0 ? "" : line.slice(separator + 1)).replace(
      /^ /u,
      "",
    );
    if (field === "data") data.push(value);
    if (field === "event") event = value;
    if (field === "id" && !value.includes("\0")) id = value;
    if (field === "retry" && /^\d+$/u.test(value)) retry = Number(value);
    return undefined;
  };
  try {
    while (true) {
      const chunk = await reader.read();
      if (chunk.done) break;
      buffer += decoder.decode(chunk.value, { stream: true });
      while (buffer.includes("\n")) {
        const newline = buffer.indexOf("\n");
        const raw = buffer.slice(0, newline);
        buffer = buffer.slice(newline + 1);
        const frame = consume(raw.endsWith("\r") ? raw.slice(0, -1) : raw);
        if (frame) yield frame;
      }
    }
  } finally {
    await reader.cancel().catch(() => undefined);
    reader.releaseLock();
  }
}

const terminal = (run: Run) =>
  ["completed", "failed", "cancelled", "timed_out"].includes(run.state);

async function getRun(runId: string): Promise<Run> {
  return request<Run>(`/projects/${projectId}/runs/${runId}`);
}

async function waitForRun(runId: string): Promise<Run> {
  let run = await getRun(runId);
  let reconnectDelay = 1_000;
  while (!terminal(run)) {
    const headers = new Headers({
      authorization: `Bearer ${apiKey}`,
      accept: "text/event-stream",
    });
    if (state.lastEventId) headers.set("last-event-id", state.lastEventId);
    const response = await fetch(`${origin}/v1/projects/${projectId}/events`, {
      headers,
      signal: shutdown.signal,
    });
    if (!response.ok || !response.body) {
      throw new Error(`SSE connection failed with ${response.status}`);
    }
    for await (const frame of parseSse(response.body)) {
      JSON.parse(frame.data); // Apply or validate the public event first.
      if (frame.retry !== undefined) reconnectDelay = frame.retry;
      if (frame.id) {
        state = { ...state, lastEventId: frame.id };
        await saveState();
      }
      run = await getRun(runId);
      if (terminal(run)) return run;
    }
    await new Promise<void>((resolveDelay, reject) => {
      const timer = setTimeout(resolveDelay, reconnectDelay);
      shutdown.signal.addEventListener(
        "abort",
        () => {
          clearTimeout(timer);
          reject(new DOMException("Aborted", "AbortError"));
        },
        { once: true },
      );
    });
  }
  return run;
}

async function main(): Promise<void> {
  if (!state.sessionId || !state.runId) {
    const source = await readFile(contextPath);
    if (source.length > 90_000 || source.includes(0)) {
      throw new Error("Trusted context must be small UTF-8 text, not binary");
    }
    const text = new TextDecoder("utf-8", { fatal: true }).decode(source);
    const initialMessage = [
      "Review this operator-selected, trusted text excerpt.",
      `Path label: ${process.env.OAO_CONTEXT_LABEL ?? "selected-source"}`,
      "```text",
      text,
      "```",
    ].join("\n");
    if (initialMessage.length > 100_000) {
      throw new Error("The complete initialMessage exceeds 100,000 characters");
    }
    state.createSessionKey ??= required("OAO_CREATE_SESSION_KEY");
    await saveState();
    const created = await request<CreatedSession>(
      `/projects/${projectId}/sessions`,
      {
        method: "POST",
        body: JSON.stringify({
          agentId,
          title: "Trusted code excerpt review",
          initialMessage,
        }),
      },
      state.createSessionKey,
    );
    state = {
      ...state,
      sessionId: created.id,
      runId: created.latestRunId,
    };
    await saveState();
  }

  let run = await waitForRun(state.runId!);
  console.log(`Initial run ${run.id}: ${run.state}`);

  const followUp = process.env.OAO_FOLLOW_UP;
  if (followUp && !state.followUpRunId) {
    if (!terminal(run)) throw new Error("Latest run has not settled");
    state.followUpKey ??= required("OAO_FOLLOW_UP_KEY");
    await saveState();
    const next = await request<Run>(
      `/projects/${projectId}/sessions/${state.sessionId}/runs`,
      { method: "POST", body: JSON.stringify({ message: followUp }) },
      state.followUpKey,
    );
    state = { ...state, followUpRunId: next.id };
    await saveState();
  }

  if (state.followUpRunId) {
    run = await waitForRun(state.followUpRunId);
    console.log(`Follow-up run ${run.id}: ${run.state}`);
  }

  const session = await request<unknown>(
    `/projects/${projectId}/sessions/${state.sessionId}`,
  );
  console.log(JSON.stringify(session, null, 2));
}

try {
  await main();
} catch (error) {
  if (!shutdown.signal.aborted) throw error;
  console.log("Shutdown requested; restart with the same state file and keys.");
} finally {
  shutdown.abort();
  process.removeListener("SIGINT", stop);
  process.removeListener("SIGTERM", stop);
}
````

Choose and persist idempotency keys outside the source tree, then run:

```bash theme={null}
export OAO_AGENT_ID='your-agent-id-with-no-caller-tools'
export OAO_TRUSTED_TEXT_FILE=/absolute/path/to/reviewed/file.ts
export OAO_CREATE_SESSION_KEY='repo-review-2026-08-20-001'
export OAO_FOLLOW_UP='List the three highest-risk boundaries.'
export OAO_FOLLOW_UP_KEY='repo-review-2026-08-20-001-follow-up-1'
node --experimental-strip-types codebase-inline.ts
```

For an external caller-tool worker, use the same HTTP paths as the SDK methods:

| SDK method | HTTP path |
| - | - |
| `listToolCalls(projectId, { runId })` | `GET /v1/projects/{projectId}/tool-calls?runId={runId}` |
| `claimToolCall` | `POST /v1/projects/{projectId}/tool-calls/{id}/claim` |
| `renewToolCall` | `POST /v1/projects/{projectId}/tool-calls/{id}/renew` |
| `releaseToolCall` | `POST /v1/projects/{projectId}/tool-calls/{id}/release` |
| `submitToolResult` | `POST /v1/projects/{projectId}/tool-calls/{id}/result` |

Keep the same repository executor, result envelope, lease heartbeat, fence, and
idempotency rules. Do not weaken the boundary because the worker lives in a
different repository.

## Failures and recovery

| Failure | Correct response |
| - | - |
| `401 unauthenticated` | Replace or correct the server credential. Never move it into the URL. |
| `403 forbidden` | Check that the key belongs to the path's project and has the exact required scope. |
| `409 conflict` when sending a turn | Re-read the session. Submit only after its latest run is `completed`, `failed`, `cancelled`, or `timed_out`. |
| `409 idempotency_conflict` | The same mutation key was reused with a different body. Recover the original request; do not hide the conflict with a new key. |
| Claim is not available | Another principal holds an unexpired lease, or the work is no longer pending. Re-list durable tool calls. |
| Stale tool fence | Discard the result. Re-list and reclaim; never submit under the old fence. |
| Network failure after a write | Retry the identical route and body with the same persisted idempotency key. |
| SSE disconnect or normal close | Reconnect with the last event ID saved after successful application. |
| SSE stalls or drops behind a proxy | Disable buffering and compression for `text/event-stream`. Allow at least 15 seconds between reads; OAO keeps quiet streams alive. |
| Unknown event kind | Persist nothing yet, re-fetch the affected durable resource, apply it idempotently, then advance the cursor. |
| Worker crash during a tool | The lease expires. A restarted worker re-lists `caller_pending` and `caller_claimed` work and claims it with a new fence. |

Run failures are durable terminal outcomes, not transport exceptions. Read the
run, session transcript, safe timeline, and public request ID before deciding
whether to submit a new turn. `POST /runs/{runId}/resume` exists for a settled
run and creates another run in its session, but it does not continue an
unsettled execution or bypass the one-active-run invariant.

<Note>
  Authorized session views include provider thinking and tool arguments/results.
  Public events and application logs still omit that content, credentials,
  authorization headers, and provider reasoning signatures.
</Note>


This documentation is built and hosted on [Mintlify](https://mintlify.com), a developer documentation platform.