Files
PageIndex/pageindex/agent_tools.py
T
Ray 4e41acdc68 feat: agent tools and local chat for the PageIndex SDK (v0.2.10) (#396)
* feat: agent tools — the cloud MCP tool contract on the client

Four new client methods make PageIndex documents available to agent
frameworks, in both modes, with the mode decided solely by the client
constructor:

- agent_tools(): plain functions (browse_documents, get_document,
  get_document_structure, get_page_content) matching the PageIndex cloud
  MCP server's tools/list — same names, schemas, descriptions, and JSON
  response envelopes — so agent prompts port unchanged between the cloud
  MCP connection and these in-process tools. Tools never raise; errors
  come back in the same envelope. remove_document ships behind
  include_management=False.
- as_openai_tools(): the same tools wrapped for the OpenAI Agents SDK.
- as_claude_mcp(): one mcp_servers entry for the Claude Agent SDK —
  cloud clients get the remote MCP config (the framework connects to
  api.pageindex.ai/mcp and discovers the full cloud tool set), local
  clients get an in-process SDK MCP server.
- agent_instructions(doc_id=None): orchestration guidance for the
  agent's system prompt; doc_id (same shape as chat_completions) appends
  the target documents.

submit_document() gains wait=True: poll get_document status until
completed, raise on failed or after 30 minutes — the manual polling loop
every cloud caller writes today spins forever on a failed document.

Neither framework becomes a dependency: imports happen at call time with
actionable errors, and pageindex[openai] / pageindex[claude] extras are
floor-only pins. tests/data/cloud_mcp_contract.json freezes the tool
contract; a parity test guards against drift. 36 new tests (95 total),
plus a live OpenAI Agents SDK run over a seeded local store verifying
the structure-first navigation flow end to end.

* fix: agent tools review — next_steps order, resolve caching, error semantics

- Large-doc next_steps now says structure-first, consistent with tool
  descriptions and agent instructions
- _remove_document fetches document list once instead of per-name
- call_tool returns error envelope for unknown names instead of raising
- _not_ready_error timed_out flag reflects actual wait outcome
- openai_agents.py docstring corrected to match default (FunctionTools)
- Removed unused ModelSettings import from demo

* fix: agent tools review 2 — bridge thread safety, browse paging, metadata merge

- McpBridge reads session/protocol headers under the lock (now RLock:
  _ensure_initialized posts while holding it). openai-agents runs sync
  tools on threads and executes parallel tool calls concurrently, so
  bridge functions genuinely race; a torn read sent a new session id
  with a stale protocol header. Measured: one session expiry under 8
  threads cost 4 initializations before, minimal 2 after.
- Session-expiry retry also resets the negotiated protocol version, so
  the re-handshake carries no stale MCP-Protocol-Version header.
- browse_documents time sort pages list_documents natively instead of
  fetching the whole library to slice one window (relevance still needs
  the full list for scoring).
- _await_completion: a status refetch that nulls out metadata no longer
  clobbers the listing's copy (setdefault was a no-op on existing None).
- Structure tool reads the raw stored tree via a named LocalAPI
  raw_tree() seam instead of reaching into _api._store internals; drop
  the redundant deepcopy before _format_structure (store re-reads from
  disk, formatting builds fresh containers).
- Shared pageindex/_version.py replaces _sdk_version duplicated in
  mcp_bridge and the Claude integration.

Left as-is after source verification against the cloud MCP: first-page
budget bypass, pageNum falsy-zero, and the page-gap fallback text are
letter-for-letter cloud behavior — parity wins over local repair.

* fix: agent tools review 3 — page-span cap, duplicate names, wait resilience, contract drift

- _parse_page_spec bounds the requested span arithmetically (10k pages)
  before materializing it; pages="1-1000000000" previously expanded to a
  billion integers inside the caller's process.
- Local submit_document uniquifies document names the way the cloud
  upload does (taken name -> _1.._99, then reject with the cloud's own
  message). Same-name duplicates broke name-addressed tools: resolution
  always picks the newest, so older duplicates were unreachable.
- agent_instructions(doc_id=...) now fails loud when the pinned doc's
  name is shadowed by a newer same-name document (legacy stores predate
  the rename) — it previews resolution with the same _resolve_document
  the tools use, so the check cannot drift from actual behavior.
- submit_document(wait=True) tolerates transient network errors, not
  just API errors; a dropped connection at minute 25 of a 30-minute
  wait no longer kills it. Third strike wraps into PageIndexAPIError
  per the documented contract.
- The live contract-parity test compares full per-param schemas, not
  just names and descriptions. It immediately caught real drift the
  shallow check had been passing: the server now emits nullables as
  anyOf unions and stamps MAX_SAFE_INTEGER maxima on offset/part.
  Contract and snapshot updated to the served wire form; _annotation_for
  learned anyOf so bridge signatures stay Optional[str] instead of
  degrading to Any.

Adjudicated, not changed: the allowed_tools wildcard example stays
(docstring advice covers scoping; Ray's call), and raw-length response
accounting stays (letter-for-letter cloud behavior, parity wins).

* feat: surface the stored document name from submit_document

Compute PR #558 makes /doc/ return {"doc_id", "name"} carrying the
post-dedup-rename name. Mirror it end to end: local submit returns the
stored name, the client warns when it differs from the uploaded file
name (read via .get so older cloud servers stay compatible), the local
name-exhaustion check runs before indexing instead of after the LLM
spend, and the demo caches doc_id in a file instead of name-matching —
a renamed document made the name lookup re-index on every run.

* fix: add missing page_list kwarg in duplicate-name test mock

* revert: keep README.md unchanged from main — SDK section deferred

* feat: serve cloud agent instructions live from the MCP server

The cloud MCP server publishes its agent instructions in the initialize
result, adapted to each key's tool set. agent_instructions() previously
returned the SDK's local-subset text in both modes — a silently forked
copy that lacks the guidance for cloud-only tools (search_documents
escalation, folders, images) and drifts as the server's prompt evolves.

Cloud clients now serve the server's live instructions, captured from
the initialize handshake on a per-client bridge shared with
agent_tools() (one session, no extra request). An empty server response
raises instead of silently substituting the subset text — same posture
as the annotation-regression guard. The local constant stays as the
honest subset for the in-process tools, with its provenance noted and a
consistency test that every tool it names exists in the local registry.

* fix: local relevance sort answers honestly instead of imitating

sort="relevance" is cloud-side semantic ranking; the local substring
imitation could satisfy the letter of the interface while silently
missing semantically relevant documents. Per the honest-subset rule
(same treatment as folders), local now returns the "not available
here" envelope for sort="relevance" or a stray query, and the local
instructions steer discovery through name/description matching plus
full-library paging instead of prescribing a capability that does not
exist here. The tool schema keeps the cloud contract verbatim, like
folder_id: honesty lives in the runtime answer, not a forked contract.

* docs: note the cloud+Claude instructions duplication trade-off in as_claude_mcp

* fix: unsupported-capability envelopes say local-mode-yet, point to cloud

"Not available here" read as a broken feature; the honest framing is
that folders and semantic ranking exist on PageIndex cloud and are not
in local mode yet. Both envelopes now say so and name the cloud client
in next_steps, so agents relay an accurate story to the user.

* fix: local tool descriptions pre-announce cloud-only capabilities

The cloud-verbatim browse_documents description invites
sort="relevance" and folder drilling, so a local agent's first semantic
search attempt was a guaranteed dead end discovered only from the
runtime error envelope. Local registration now appends a LOCAL MODE
note to the description — the agent learns what is cloud-only before
calling; the runtime envelope stays as the backstop for prompts that
ignore descriptions. The cloud-facing contract stays byte-verbatim.

* refactor: localized tool guidance replaces the appended LOCAL MODE note

Appending a retraction to the cloud-verbatim description left the model
parsing an instruction and its negation — and kept the cloud text
recommending search_documents and get_folder_structure, tools that are
not registered locally (get_page_content likewise pointed at
get_document_image). Guidance now adapts to the local surface the way
AGENT_INSTRUCTIONS already does: schema structure stays byte-identical
to the contract (mechanically asserted by a strip-descriptions test),
while local description strings teach only what works here and point to
PageIndex cloud for the rest. A dead-reference test forbids local
guidance from naming tools outside the local registry, so a contract
refresh that reintroduces a cloud-only reference fails loudly.

* feat: hide cloud-only parameters from the local tool surface

folder_id, sort, query, and recursive were exposed locally with
localized "cloud-only" descriptions, leaving the dead-end calls
expressible and discovered at runtime. Schema constraints beat
guidance: the local surface now serves the contract minus these
parameters, so strict-schema frameworks make the calls inexpressible
and a prompt that insists on sort="relevance" degrades to the bare
call (the correct local behavior) instead of an error round-trip.

The implementations still accept the hidden parameters and answer with
the guided "works on PageIndex cloud" envelope — the backstop for
direct call_tool callers and hosts without schema enforcement.
wait_for_completion stays: seeded or torn stores can hold documents
that are genuinely not completed. The structural guard now asserts the
local schema equals the contract minus the documented hidden set,
descriptions aside.

* fix: incremental-review findings — bridge cache, guards, envelope drift

Three independent review passes over the agent-instructions increment
surfaced six fixes:

- The per-client bridge moved off the instance into a weak-keyed,
  lock-guarded module cache: cloud clients stay picklable
  (threading.RLock no longer rides on the client) and concurrent first
  calls can no longer construct duplicate bridges/sessions.
- Blank or non-string initialize.instructions now hit the same honest
  error as a missing one — a whitespace-only or structured value could
  previously become the system prompt (or crash the doc_id append with
  a raw TypeError).
- The invalid-sort envelope no longer prescribes sort="relevance" — the
  one error text that still taught the cloud-only value it would then
  reject.
- "Page through the rest of the library" is emitted only when has_more
  is true; a fully-listed library no longer instructs a pointless call.
- The mandatory full-library paging step now says limit: 50 — 6 calls
  instead of 30 on a 300-document library.
- Docstrings and comments rescoped to what is actually true: the
  never-raise contract covers invocations the signatures accept
  (unknown params fail at the Python boundary; call_tool answers them
  with the guided envelope), recursive is accepted as the identity
  rather than errored, lenient framework arg models drop hidden params
  pre-call, and the module header no longer claims full schema parity.
  The capability-phrase guard now covers every local docstring, not
  just browse_documents.

* chore: keep the demo's doc_id cache file out of the repo

* test: live envelope field-parity guard against cloud response drift

The frozen contract guards tools/list, but the response envelopes the
local tools emit were hand-built to mirror the cloud's and had no drift
detector. A key-gated live test now asserts every field local emits
exists in the live cloud response for the analogous call (top-level
keys, next_steps, document entries, structure nodes, content entries).
Guidance wording is deliberately localized and not compared. Verified
green against the live server: local and cloud field structures
currently match exactly.

* feat: local chat — three protocol surfaces over the agent tools (v0.2.10)

Local mode gains managed document QA: an agent over the #393 local tool
set, reachable through three wire protocols, each 1:1 with the backend
and with no translation layer.

- chat_completions(): standard chat.completions semantics on any
  OpenAI-compatible backend (openai-agents engine). Final answer only,
  cross-turn aggregated usage, streaming as text pieces or chunk dicts
  (the existing cloud signature, now implemented locally; model and
  max_turns are local-only additions).
- responses(): the agentic surface — OpenAI Responses format, the tool
  process is standard output items, streaming forwards native events
  (tool outputs emitted as response.output_item.done, the way the
  platform streams its own server-side tools). Round-tripping output
  into the next input keeps provider prompt-cache prefix continuity and
  the agent's memory — live-verified: the follow-up call answered from
  round-tripped tool output with zero new tool calls.
- messages(): Anthropic-native via the SDK's own tool runner (new
  pageindex[anthropic] extra, floor 0.68.0 verified for
  tool_runner/beta_tool(input_schema)). tool_use/tool_result round-trip
  is the format's native behavior; the envelope is the final message
  with aggregated usage plus the full new-turn sequence; the managed
  system blocks carry cache_control breakpoints.

Shared skeleton: thin chat header + the local AGENT_INSTRUCTIONS
(caller system content is appended, not rejected), the doc_id targeting
block as a leading context item (factored out of
build_agent_instructions), read-only toolset, structural-only
validation (no arbitrary caps — backend limits govern), sampling params
passed through, per-run tracing disabled, enable_citations rejected as
cloud-only. Design basis is industry-standard formats rather than the
cloud chat endpoint; responses()/messages() raise on cloud clients
until the cloud converges.

Tests run the real engines against scripted backends (a Model fake for
openai-agents, a mock HTTP transport under the real anthropic SDK) with
real tool execution against a seeded store, including the round-trip
prefix-extension assertions on both engines.

* fix: local-chat review findings — truncation, serialization, streams

Three independent review passes (bug scan, claims-vs-code, adversarial
runtime probes) over the local-chat increment; every fix below was
reproduced before being fixed.

messages():
- A max_turns cut no longer duplicates the final assistant turn: the
  runner has already appended it when iterations exhaust, so the
  round-trip history carried a duplicate tool_use id and ended on an
  unanswered tool_use — a guaranteed 400 on continuation. The append
  now keys on stop_reason, and truncation reads natively as
  stop_reason: "tool_use" with a continuable history.
- The envelope is JSON-serializable end to end: runner-stored turns
  carry pydantic content blocks; everything is dumped to plain dicts,
  excluding SDK-internal __api_exclude__ fields (parsed_output) that
  the API rejects on round-trip.
- Bounded by default (max_iterations 10, like the OpenAI surfaces);
  usage aggregation now preserves the final turn's native fields and
  sums the token counters None-safely; empty caller system strings are
  skipped; non-dict message entries and bad doc_id types raise
  PageIndexAPIError; anthropic < 0.68 gets an actionable version error;
  the doc block no longer spends a cache_control breakpoint.

chat_completions()/responses():
- MaxTurnsExceeded wraps into PageIndexAPIError on all four run paths.
- responses(stream=True) is one logical response: per-turn backend
  lifecycle events are collapsed (a canonical consumer previously
  stopped at turn 1's response.completed and never saw the answer),
  sequence numbers are reassigned monotonically, and the synthesized
  tool-output event carries output_index/sequence_number.
- The responses envelope carries the real request surface
  (instructions, the actual function tool definitions, tool_choice,
  parallel_tool_calls, error/incomplete_details).
- RunConfig(group_id) pins a stable prompt_cache_key: openai-agents
  otherwise stamps each run with a fresh key, tagging round-tripped
  prefixes as different cache groups and defeating the feature the
  round-trip exists for.
- Abandoning a stream now cancels the run: a watchdog task lets the
  cancellation land even while the pump awaits the backend, and the
  per-call AsyncOpenAI client is closed before its loop ends (fixes
  "Task exception was never retrieved" noise). The opening role chunk
  is emitted even for empty outputs; empty responses() input and
  enable_citations-before-extra ordering fixed.

Docs rescoped to what is true: finish_reason/status reflect loop
completion on the OpenAI surfaces (the engine does not surface per-turn
backend reasons); chat streaming yields visible narration including
pre-tool text; messages(stream=True) forwards the Anthropic SDK's
native event objects (not wire-verbatim); the doc block is a leading
conversation item on OpenAI surfaces and a system block on messages().

Tests: 25 in the file (11 new), with per-extra skip sections so a
machine with only one framework still covers the other surface;
without-frameworks matrix re-verified; live smoke re-run green with a
clean exit.

* feat: as_anthropic_tools — Anthropic tool-runner export, both modes

Fills the last cell of the agent-connection matrix: users driving their
own anthropic tool_runner loop get runnable tools directly. Cloud wraps
the live MCP tool set with input schemas passing through verbatim (MCP
inputSchema is the Messages API schema shape); local exposes the same
set messages() runs internally. The beta_tool wrapping moves from
local_chat into integrations/anthropic_sdk.py, parallel to
openai_agents.py, and messages() now consumes the shared builder.
agent_tools grows _bridge_invoker/_read_only_tools so the plain-function
and beta_tool cloud paths share invocation containment and the
read-only gate.

* fix: as_anthropic_tools review findings — async flavor, schema isolation

Adversarial + best-practice review of 4590dd8 (three independent passes)
surfaced two holes. The export was sync-only: AsyncAnthropic's runner
accepts only BetaAsyncFunctionTool and splices anything else into the
request body unserialized, so the first call died with an opaque
TypeError — asynchronous=True now builds beta_async_tool runnables
(present since the 0.68.0 floor) that run the blocking bridge/store call
in a worker thread, keeping I/O off the caller's event loop. And
beta_tool stores input_schema by reference, so cloud tools aliased the
bridge's cached metas while the local path deep-copied — the builder now
copies, and the passthrough test asserts equal-but-not-aliased so it can
no longer compare an object with itself. Docstring fixes from the same
round: the MCP-connector pointer now carries the full live-verified
shape (authorization_token was missing — following it literally gave a
401), and the manual messages.create loop's to_dict() serialization is
documented. Tests pin the runnable flavor both ways (isinstance), which
existing tests could not distinguish.

* docs: doc_id is per-call table-setting — keep it identical across a conversation

The targeting block doc_id adds is re-set on every call and sits in the
cached prompt prefix, so a round-trip that drops (or changes) doc_id
silently diverges the prefix and loses the cache continuation. State the
rule on all three chat surfaces' doc_id docs, and pin it with a prefix
test that passes the same doc_id on both calls.

* feat: every chat surface takes a bare query string

query + doc_id is the minimal PageIndex contract, so it now works
uniformly: chat_completions and messages accept a plain string (one
user message), as responses always did per its wire format. The wrap
is input sugar at the SDK surface, not a translation layer — the
outgoing wire is unchanged, and managed agent surfaces taking strings
is the ecosystem convention (Runner.run, claude_agent_sdk.query).
Cloud chat_completions gains the same acceptance; blank strings raise
on every path.

* feat: messages() defaults max_tokens to 4096

The Messages API requires a per-turn output budget on the wire, but
that is table-setting, not a PageIndex-layer user obligation — the
simple call is now a question + model + doc_id. The knob stays
overridable (passthrough intact); model stays required because no
cross-vendor default is honest to guess.

* fix: raise messages() max_tokens default to 8192

max_tokens is a cap, not consumption, so the default should be the
highest universally safe value: 4096 could truncate long-form answers
(whole-document summaries), while 8192 is the output ceiling every
non-EOL Claude model accepts and stays under the SDK's non-streaming
long-request threshold.

* fix: restore per-extra skip markers the string-input tests displaced

Inserting tests above decorated ones absorbed their @needs_agents
markers, so two tests ran (and failed) in the without-frameworks CI
job. Both simulated-bare and full runs are green again.

* fix: close 17 findings from the v0.2.10 max review

Tool layer:
- anthropic adapter: failed tool calls raise ToolError so the runner
  emits tool_result is_error:true; McpBridge.call_tool returns
  (text, is_error) and surfaces the server's MCP isError marking
- as_openai_tools builds FunctionTool with the contract/server schema
  verbatim (strict off) — function_tool() regenerated schemas from
  signatures, dropping items/enum/pattern/bounds and aborting the whole
  list on object-typed params; shared _tool_specs() feeds both adapters
- remove_document validates every name before deleting anything; call_tool
  classifies only bind-time TypeErrors as INVALID_INPUT
- unknown-tool envelope formatted with _dumps like every other envelope

Local chat:
- doc_id is enforced at the tool layer (allowlist threaded through
  call_tool and the adapters), not just prompted; the shadow check runs
  inside the scope
- _openai_model routes litellm/ and provider/ paths via LitellmModel and
  strips openai/ — the normalized retrieve_model 404'd as a raw wire name
- responses() reports the backend's real terminal status (recorded at the
  transport client; the framework discards Response.status) and wraps
  framework exceptions in PageIndexAPIError
- chat_completions streaming yields its opening chunk inside try, so an
  abandoned iterator still cancels the run and closes the backend
- prompt-cache group_id is per-conversation (model+instructions+first
  item) instead of one global constant pooling every user
- messages() max_tokens default resolves per model (claude-3 caps at 4096)

Packaging / surface:
- __init__ registers the 0.2.10 modules in _SUBMODULES; unknown names
  raise AttributeError instead of eagerly importing page_index_classic
- anthropic floor 0.84.0: first release with ToolError whose runner also
  executes the final turn's tools on a max_iterations cut
- client docstrings caught up with local chat landing

Claude Agent SDK gate:
- claude_allowed_tools(mcp_servers) derives mcp__<key>__<tool> entries
  from the caller's own registration map (live server annotations on
  cloud, the contract locally) — no name is ever spelled twice
- claude_agent_config() bundles the three slots as one-call sugar over
  the explicit form

Examples:
- demo runs against cloud again (getattr for local-only attrs) and finds
  an existing indexed copy by name before re-indexing

Tests: monkeypatches replace the consuming module's binding instead of
mutating the shared time/requests modules; 185 -> 211.

* feat: one-call config bundles for every bring-your-own-framework surface

claude_agent_config() gets two symmetric siblings, so each framework's
front door is a single splat over the same explicit primitives:

- openai_agent_config(): Agent(**...) kwargs — instructions, tools, and
  the local retrieve_model (cloud omits model for the framework default)
- anthropic_runner_config(): tool_runner(**...) kwargs — system, tools,
  and the messages() defaults (per-model max_tokens, 10-iteration bound);
  only the user's messages remain

Bundles stay pure sugar: doc_id rides agent_instructions, no extra
semantics over the explicit form, docstrings point both ways. The demo
agent shrinks to Agent(**client.openai_agent_config(doc_id=...)).

Construction is pinned against the real frameworks in tests (Agent and
tool_runner both built offline), so an upstream kwargs rename fails
loudly; 211 -> 215 tests.

* fix: three more review findings — partial-read reporting, reply correlation, output_index axis

- get_page_content: the summary is additive, not either/or — a call that
  both truncates for size and has out-of-range pages reported only the
  latter, telling the agent every in-range page was returned (#2)
- McpBridge._extract_result: strict request-id correlation only; the
  eager fallback could hand back a stale or mis-correlated JSON-RPC
  message as this call's reply (#16)
- responses() streaming: output_index now addresses the logical
  response.output — backend per-turn indexes are re-based past prior
  turns' items and the SDK-injected tool outputs take the next slot on
  that axis, instead of reusing the event-sequence counter (#15)

215 -> 217 tests.

* feat: gate the config-handoff surfaces by the read-only MCP endpoint

pageindex-chat#448 adds /mcp?tools=read — the server registers only
readOnlyHint-annotated tools — so the URL itself becomes the gate for
every surface that hands a config to a third party:

- as_claude_mcp: include_management now picks the endpoint on cloud;
  the parameter is real in both modes
- as_openai_tools(hosted=True): OpenAI connects to the read-only
  endpoint by default and require_approval simplifies to "never" — the
  approval-flow middle ground becomes hard absence, matching every
  other surface's default
- claude_allowed_tools() retired before ever shipping: with the server
  gated, allowed_tools degenerates to whole-server pre-approval, which
  claude_agent_config emits as the constant ["mcp__<name>"] — no
  setup-time bridge round-trip remains
- in-process surfaces (agent_tools, as_openai_tools, as_anthropic_tools
  over the bridge) keep bare /mcp + client-side annotation filtering:
  they materialize tools locally and hand no URL to anyone

Release ordering: 0.2.10 must ship after pageindex-chat#448 deploys —
an older server ignores unknown query params and would silently serve
the full set behind a URL that promises read-only.

* fix: chat_completions wraps framework exceptions like responses()

The AgentsException -> PageIndexAPIError wrap from the responses() fix
covered only that surface; a backend stream dying without a terminal
event (or any engine failure) still escaped chat_completions as a raw
openai-agents exception type on both its paths.

* fix: same-name documents in different folders no longer refuse doc_id targeting

The shadow check in doc_targeting_block compared names across the whole
library, so agent_instructions(doc_id=...) hard-raised for a legal cloud
layout — one file name in two folders — with advice (rename/remove) that
contradicts the contract, whose folder_id parameter exists precisely to
disambiguate this case.

Shadowing is now judged per folder: only a newer same-name document in
the SAME folder makes the name unreachable and raises. A same-name
document in another folder serves the call, and the targeting block adds
a directive to pass folder_id on every tool call — dropping the raise
alone would have traded a loud refusal for the agent silently reading
the newer document.

Local mode (folderId always None) and the scoped chat path (allowlist
resolution, fixed with the doc_id enforcement) are behaviorally
unchanged.

* Revert "fix: same-name documents in different folders no longer refuse doc_id targeting"

This reverts commit 9d16dbf6, whose premise collapsed on verification
against the cloud upload paths. Review finding #10 inferred from the
folder_id tool description that one file name in two folders is a legal
cloud layout; both upload paths actually dedup names per USER SPACE with
no folder dimension — chat's getSignedUploadUrl queries fileName +
sourceName + mode + owner (file-access.service.ts), and compute's
get_upload_url probes the S3 key (user, source, file_name) — so own
documents cannot share a name across any folders. The only legitimate
same-name source is shared mounts (shared-with-me/following), which the
api-proxy surface this SDK talks to never carries.

Cloud and local therefore share one invariant — names unique per space,
server-enforced — and the original global shadow check was the right
shape: a duplicate is an anomaly worth refusing loudly, not a layout to
accommodate with per-folder adjudication and conditional prompt notes.
The invariant is now stated in doc_targeting_block's docstring so the
finding does not get re-raised.

* docs: state the name-uniqueness invariant in library terms

* docs: trim doc_targeting_block docstring to the contract

* fix: compress out-of-range page lists in get_page_content

Two message strings enumerated every out-of-range page number one by one
while the payload fields beside them already used _format_page_spec.
Against a 2-page document, pages="3-10000" produced a 59,310-character
error whose own requested_pages field expressed the identical set as
"3-10000"; the mixed case pages="1-10000" produced 59,479. Both now
render through the helper: 431 and 600 characters.

This was inherited behaviour, not a local slip — the cloud MCP server
enumerated at the same two sites, so local reproduced it verbatim. The
cloud fixed it first (pageindex-chat #449), and this follows to keep the
strings byte-identical; the error message now matches the served one
character for character. A differential run of the two compressors over
94 inputs (empty, single, unsorted, duplicated, 10k spans, 80 random)
agrees on every one, separator included.

The new test pins all three shapes, including the non-contiguous case
("5,9" must not collapse into a range) that the compressor had no direct
coverage for.

* docs: messages() marks only the managed prefix with cache_control

The method docstring claimed the doc targeting block carries a
cache_control breakpoint too, and the doc_id note called that block part
of the cached prompt prefix. _anthropic_system deliberately marks only
the stable managed prefix — the API allows four breakpoints and the
varying doc block must not consume one — and the block is appended after
the sole breakpoint, so it is never cached.

a45b554 added the block with cache_control, making the claim true when
written; daac9d2 removed it without touching the docstring, and adb2f1f
then added the "cached prompt prefix" sentence after the fact. The same
phrase at the chat_completions and responses docstrings is correct —
there the block is a leading conversation item inside the auto-cached
prefix — so only the Messages surface is reworded.

The doc_id advice itself stands: the block is per-call table-setting and
should stay identical across a conversation. Only the caching rationale
was wrong.

* docs: _stream_sync cancels on close, not on abandonment

The docstring promised that closing "or abandoning" the iterator cancels
the run. Abandoning only works when refcounting collects the generator:
a caller that breaks out of the loop while keeping the reference never
runs the finally that sets the cancel event, so the pump thread stays
parked on the full queue and the backend client is never released.

Closing is correct and is what the dedicated test exercises. Narrowing
the promise to the behaviour the code actually provides is the honest
fix; a watchdog or finalizer would be machinery bought for a shape the
sync surface is not meant to serve, and the async client planned for
0.2.11 gets native task cancellation instead.

* fix: raise the openai-agents floor to 0.14.0

_conversation_group_id feeds RunConfig.group_id into OpenAI's
prompt_cache_key so a round-tripped prefix stays in one cache group.
That wiring first appears in openai-agents 0.14.0: 0.8.0 through 0.13.x
have no prompt_cache_key at all, and group_id there is a tracing group
id only — inert, since tracing is disabled on the line above. An install
resolving to the declared floor lost the cache continuity that the
responses() docstring sells, silently and with no test able to catch it.

The old floor's rationale (0.8.0 offloads sync tools to a thread) is
subsumed by the new one. Every symbol the package imports predates
0.14.0, so nothing else constrains the bound.

* fix: enforce doc_id at the tool layer in the framework config helpers

openai_agent_config / anthropic_runner_config / claude_agent_config
accepted doc_id but built unscoped tools, so the parameter that is a
structural allowlist on chat_completions() was prompt-only advice here —
the agent could read every document in the store regardless.

- as_openai_tools / as_anthropic_tools / as_claude_mcp take a doc_id
  tail parameter and thread it to the existing _allowed_ids channel;
  the config helpers pass it through in local mode
- cloud config helpers keep prompt-level targeting (tool scoping is
  server-side there, documented); explicit as_*(doc_id=...) raises on
  cloud instead of silently dropping the allowlist — including the
  hosted branch, which returned before _tool_specs' existing guard
- _require_local_scope consolidates the cloud rejection that was
  inlined in _tool_specs
- doc_id=[] is an empty allowlist, not "unscoped": dropped the
  `or None` at the three local chat surfaces

* fix: two chat findings — final-turn append and cache-key seeding

run_messages keyed its re-append guard on stop_reason, but the anthropic
runner executes tools whenever the turn's content carries tool_use
blocks (refusal excepted) — a max_tokens turn with complete tool_use
blocks was already appended by the runner, so the guard re-appended it,
duplicating tool_use ids and 400ing the documented verbatim
continuation. The guard now checks whether final's tool_use ids already
sit in the appended history; unexecuted tool_use blocks (refusal turns)
are stripped from the appendable history, as the SDK itself does when
rebuilding params around an unresulted turn.

_conversation_group_id seeded on items[0], which is the doc-targeting
block whenever doc_id is set — byte-identical across every conversation
about a document, so all of them pooled under one prompt_cache_key and
evicted each other's prefixes. Seed on the conversation's own first
item instead: continuations keep their key, unrelated conversations
never share one.

Also drop the dead pytestmark_openai assignment (pytest's magic name is
pytestmark; the section gate it implied never existed).

* fix: six review findings — pagination, compat, and containment

- _all_documents advances by what actually arrived and treats `total`
  as an optimization: absent/null totals and short pages silently
  truncated the library behind every name resolution
- _make_bridge_function survives description: null (the parallel
  _tool_specs path already did)
- as_openai_tools answers a malformed argument string with the guided
  error envelope instead of raising through the caller's whole run
- the pre-0.2.10 package attributes (ConfigLoader, count_tokens, ...)
  resolve again: main's underscore-guarded fallthrough is restored —
  dunder probes stay lazy, a non-underscore typo pays one classic
  import before its AttributeError
- _split_structure chunks are always lists: the structure field no
  longer changes JSON type between parts of one paginated response
- the bridge replays only session-carrying 404s (the spec's expiry
  status); 400 raises instead of re-running side effects, and the
  reset double-checks under the lock so concurrent retries cannot
  clobber a freshly re-initialized session

* fix: five secondary review findings — containment and guards

- the bridge maps content blocks individually: base64 payloads
  (image/audio) become metadata stubs instead of handing the model the
  raw blob, text blocks pass verbatim, anything else keeps the JSON
  dump (revisit if tool results become real multimodal input)
- cloud proxy annotations keep array item types (list[str], not bare
  list) so strict function calling accepts the round-trip; a type-array
  in items degrades to bare list instead of crashing the build
- run_messages raises when set_messages_params stops delivering params
  instead of silently dropping every tool turn from the envelope
- call_tool drops None-valued arguments (None ≡ omitted, the
  contract's semantics) — adapters that forward the model's nulls
  verbatim no longer trip parameter validation
- client._parse_pages bounds the span arithmetically before
  materializing it, like the tool layer: "1-999999999" raises instead
  of allocating a billion integers

* fix: three review findings — protocol honesty, model echo, containment

- responses() promised the Responses protocol ("no translation layer")
  but _openai_model ignored protocol on the LiteLLM branch: provider-
  prefixed models silently ran chat.completions under a responses-shaped
  envelope, and with no transport hook to record status (LitellmModel
  has no _client.responses) a turn truncated at the output cap reported
  status "completed". The branch now raises for protocol == "responses"
  — at agent-build time, before any backend call — naming the routes
  out: chat_completions(), messages() for Anthropic models, or
  OPENAI_BASE_URL + a bare/openai/-prefixed name for backends that
  genuinely speak /responses. Refusal, not emulation: most providers
  have no /responses endpoint to drive.
- chat_completions envelopes echoed retrieve_model verbatim, which
  carries the SDK's litellm/ routing marker after normalization — a
  name no provider catalog contains, and a different string than the
  same model passed per-call. The envelope and every streaming chunk
  now report the name the provider actually serves; routing and the
  prompt-cache group key keep the prefixed form. responses() needs no
  change (post-refusal the prefix cannot reach its envelope), and the
  user-typed openai/ prefix stays echoed as typed.
- _remove_document caught only PageIndexAPIError around the per-doc
  delete, so a bare OSError (local_store re-raises them) or a transport
  error (cloud delete_document wraps nothing) escaped mid-batch,
  discarded the entries for documents already irreversibly deleted, and
  surfaced as a generic INTERNAL_ERROR envelope inviting a retry — which
  then reports the destroyed document as not_found. The loop now catches
  Exception, keeping the per-document results the contract promises.

* fix: config bundles use the scoped shadow check their tools earned

9f67fdd made the three config helpers enforce doc_id at the tool layer
but left their instructions on doc_targeting_block's unscoped default,
so a bundle refused any doc_id whose name a newer library-wide
duplicate shadows — a raise whose message ("the tools address documents
by name and would read the newer one") had just become false: the
bundle's own tools resolve names inside the allowlist and read the
targeted document correctly. chat_completions() accepted the same
doc_id via _doc_block's scoped=True.

Each helper now computes scope = _local_doc_scope(doc_id) once and
derives both slots from it — scoped=scope is not None for the
instructions, doc_id=scope for the tools — so the check mode and the
tool allowlist come from one fact and cannot drift apart again.
build_agent_instructions grows a scoped passthrough; cloud stays on the
whole-library check (scope is None there and the tools are genuinely
unscoped), and the public agent_instructions() keeps its unscoped
default for the same reason. An in-set duplicate still raises — and in
that case the message is true on every surface that emits it.

* docs: as_openai_tools' remote-MCP note moves to the Cloud paragraph

The MCPServerStreamableHttp alternative sat in the Local: paragraph
pointing at bare {BASE_URL}/mcp — a cloud-only route (BASE_URL is the
hosted API; local has no HTTP MCP server) that as written would connect
unauthenticated to the full tool set. Now stated where it applies, in
the as_anthropic_tools connector-note form: Cloud paragraph, Bearer
auth spelled out, ?tools=read default with the drop-it escape.

* fix: nine review findings — argument coercion, scope, and honest envelopes

- call_tool coerces string booleans per the TOOL_CONTRACT schema
  ("false"/"no"/"0" read as False, not a truthy 3-minute wait) and
  survives arguments: null (json.loads("null") reaches the seam as None)
- _local_doc_scope raises on an explicitly empty doc_id on cloud: with
  no tool-layer allowlist there, dropping it silently widened an empty
  scope to the whole library
- both page-spec caps count distinct pages instead of summing parts, so
  overlapping ranges (a parent section plus its children) within the
  10k union pass again as they did in 0.2.9; the per-part arithmetic
  bound still rejects billion-page specs before materializing anything
- _remove_document deduplicates doc_names: a repeated name is one
  deletion, not a second "failed" row with an internal error string
- doc_targeting_block merges the user's metadata tags from the listing
  (local get_document keeps the 7-key cloud detail wire shape, which
  carries none) so the block delivers the metadata it promises
- _wait_until_ready folds its two raise branches into one that carries
  the doc_id: a poll that dies no longer discards the handle to an
  uploaded, billed document
- _reported_model strips both routing prefixes (litellm/ and openai/)
  and responses() now reports it too, instead of echoing a model id the
  provider never served
- _openai_model wraps AsyncOpenAI() construction so a missing backend
  credential surfaces as PageIndexAPIError like every other gate on the
  chat surfaces (and builds the client once for both protocols)
- _browse_documents advances its cursor by the rows that actually
  arrived and guards a null/absent total — the same hazards
  _all_documents already guards — and an empty window ends pagination
  instead of freezing the cursor

* fix: two chat findings — protocol terminal states, provider error types

- responses(stream=True) raised PageIndexAPIError when the backend
  ended the response with response.failed / response.incomplete:
  openai-agents yields the terminal lifecycle event, then re-raises it
  as ModelBehaviorError, so the generic AgentsException wrap
  short-circuited the emit the agen's tail was built for — its
  failed/incomplete terminal mapping was dead code against the real
  engine, and the caller lost both the partial output and the real
  status. The wrap now steps aside when the recorded terminal state is
  failed/incomplete, and the stream ends with the honest terminal
  event (committed output, real status, error/incomplete_details) —
  the backend's terminal state is a protocol event, not an engine
  failure. Non-stream was already honest for incomplete via the
  transport recorder; a failed response arrives there as an HTTP
  error, covered below. Known ceiling: the truncated final turn's
  partial text was already streamed as deltas but is not reconstructed
  into the terminal event's output (the engine commits items only on
  turn completion).
- Provider exceptions (network, auth, rate limit) leaked as raw
  openai/anthropic types through every chat surface, against the
  layer's own "never raw engine types" contract. Every engine boundary
  now wraps its vendor's base exception into PageIndexAPIError
  (chained): the four OpenAI-engine sites catch openai.OpenAIError —
  LiteLLM's exception types subclass openai's, so one handler covers
  both routing paths — and messages() catches anthropic.AnthropicError
  around the batch drive and the stream generator.

* fix: guided failure for unknown LiteLLM providers, non-object call_tool args

- _openai_model pre-checks the first path segment against
  litellm.provider_list (fail-open if the attribute ever disappears):
  a HuggingFace repo id like Qwen/Qwen2.5-7B-Instruct on an
  OpenAI-compatible server now fails at build time with the escape
  spelled out — 'openai/<id>' plus OPENAI_BASE_URL — instead of at
  request time inside LiteLLM with "LLM Provider NOT provided". The
  slash-means-provider routing convention itself is unchanged; the
  retrieve_model and chat_completions docstrings now document it where
  they promise "any OpenAI-compatible server works"
- call_tool answers a non-dict arguments value (a JSON array or scalar
  from a misbehaving caller) with the guided INVALID_INPUT envelope
  instead of raising AttributeError through the agent loop, matching
  the openai adapter's own non-object guard

* fix: wrap litellm import in PageIndexAPIError when not installed

* fix: silence CodeQL findings — merge implicit string concat, drop unused vars

* fix: two external review findings — init-notification race, SDK floor

notifications/initialized moves inside the bridge lock: a concurrent
first use could send tools/list between the handshake and the
notification, which strict MCP servers reject with a 400 the bridge
never replays. Regression test races two threads through a stalled
notification window.

claude-agent-sdk floor rises to 0.1.53 — below it, string prompts with
SDK MCP servers (the documented local-mode flow) hit invisible
registration (#597) and a deadlock (#780).

* refactor: drop the unused exc parameter from _wrap_max_turns

The parameter was dead from the moment it was introduced (daac9d2):
the body reads only max_turns, and every call site already carries the
cause via `raise ... from exc`. The signature implied the helper
inspected the engine exception, which it never did.

No behavior change — message text and __cause__ chaining verified
identical across all four call sites (chat_completions and responses,
stream and non-stream).

* fix: raise the anthropic and openai-agents floors past broken releases

Both declared floors named a version that cannot work, and CI never
caught either because it installs the latest.

anthropic >=0.84.0 -> >=0.108.0. Probed against a mock transport: on a
turn with stop_reason="refusal" carrying a tool_use block, 0.84.0,
0.92.0 and 0.100.0 all execute the tool and post the tool_result back;
0.108.0 and later stop at the refusal. test_messages_refusal_with_
tool_use_stays_appendable asserts the latter, so that test was false at
the floor. messages() is unaffected in practice (it never passes
include_management, so remove_document is not registered), but
as_anthropic_tools(include_management=True) hands it to a caller's own
runner.

openai-agents >=0.14.0 -> >=0.18.1. 0.14.0 and 0.16.0 raise pydantic
ValidationError on InputTokensDetails.cache_write_tokens before any
request reaches the transport when paired with openai 2.54.0 — and they
declare openai <3,>=2.26.0, so pip resolves exactly that pair. 0.18.1 is
clean. The 0.14.0 rationale (RunConfig.group_id -> prompt_cache_key)
still holds above the new floor.

The three extras' floor comments are cut to the binding constraint; the
reasoning lives here.

* test: cover max_turns wrapping on every chat surface

test_chat_completions_max_turns_wrapped only drove chat_completions, so
the two responses() call sites had no coverage, and no test asserted
that the engine exception survives as __cause__. Parametrized over both
surfaces and both stream modes; the non-positive max_turns rejection
splits out, since it is input validation rather than wrapping.

* fix: four review findings — envelope size honesty, contained tool errors

- _dumps drops indent=2: emission now matches _serialized_size's compact
  accounting, so the pagination budget bounds what is actually sent
  (indented parts measured under 95k but emitted ~1.8x the 100k cap)
- call_tool builds the _allowed_ids frozenset inside the guarded block:
  a non-iterable doc_id returns the INVALID_INPUT envelope instead of
  raising into the agent loop; same move for _bridge_invoker's
  arguments normalization
- next_steps strings qualify submit_document() as
  PageIndexClient.submit_document() (three sites), matching the one
  already-qualified site — it is a client method, not a registered tool
- tests: import httpx at module scope (guaranteed via the hard openai
  dependency) so agents-gated tests survive an install without the
  anthropic extra; formatting assertion follows the compact envelope
2026-08-13 07:32:11 +08:00

1673 lines
71 KiB
Python

"""Agent tools: the cloud MCP tool contract, executed against a PageIndexClient.
Tool names and the surviving input-schema structure match the PageIndex
cloud MCP server — the local surface hides the documented cloud-only
parameters — so agent prompts port across the cloud MCP connection and
this in-process layer. Only the tools that exist in every mode are
registered (no folders, search_documents, or get_document_image), and the
guidance strings (tool descriptions) adapt to the local surface the same
way the agent instructions do — they never teach capabilities that only
exist on the cloud.
Tools never raise for any invocation their signatures accept: every
outcome, including errors, is returned as the same JSON envelope the cloud
emits ({"success": true, ...} / {"error": ...}). Arguments outside a pruned
local signature fail at the Python call boundary; the call_tool path
answers them with the guided error envelope instead.
"""
from __future__ import annotations
import copy
import difflib
import inspect
import json
import re
import threading
import time
import weakref
from typing import Any, Callable, Optional
from .errors import PageIndexAPIError
TOOL_RESPONSE_CHAR_LIMIT = 100_000
STRUCTURE_FIRST_PAGE_THRESHOLD = 20
_CHAR_BUDGET = int(TOOL_RESPONSE_CHAR_LIMIT * 0.95)
_PAGES_SPEC_RE = re.compile(r"^(\d+(-\d+)?)(,\s*\d+(-\d+)?)*$")
_MAX_REQUESTED_PAGES = 10_000
_SIMILAR_NAMES_LIMIT = 3
_TOOL_WAIT_TIMEOUT = 180.0 # "up to 3 minutes", per the wait_for_completion schema
_TOOL_WAIT_INTERVAL = 5.0
_DOC_NAME_DESCRIPTION = (
'Copy the `name` field verbatim from a browse_documents() or '
'search_documents() response (case-sensitive, include extension). '
'Example: "Q3 Report.pdf". If the response shows two documents with the '
'same name, pass `folder_id` alongside to disambiguate.'
)
_FOLDER_ID_DISAMBIGUATOR_DESCRIPTION = (
'Disambiguator for same-name documents. Copy the `folder_id` from the '
'intended browse/search result; use "root" for root-level documents, or '
'"shared-with-me"/"following" for the read-only folders at the library '
'root; omit if `doc_name` is unique. Copy any folder_id verbatim from a '
'browse_documents()/get_folder_structure() response, never construct one.'
)
_WAIT_FOR_COMPLETION_DESCRIPTION = (
"If true and document is processing, automatically wait up to 3 minutes "
"until completed. Reduces repeated tool calls."
)
#: Tool names, descriptions, and parameter schemas, identical to the cloud
#: MCP server's tools/list.
TOOL_CONTRACT: dict[str, dict[str, Any]] = {
"browse_documents": {
"annotations": {"readOnlyHint": True, "openWorldHint": False},
"description": (
"Primary document retrieval tool. After orienting with "
"get_folder_structure() (when available), use this for all "
"document-related questions. The bare call returns root-level "
"sub-folders and documents; pass folder_id to drill into a "
'sub-folder level by level. Use sort="relevance" + query for '
"semantic ranking. Do NOT jump to search_documents() first — it "
"is an escalation path, only after "
'browse_documents(sort="relevance") has failed.'
),
"schema": {
"type": "object",
"properties": {
"folder_id": {
"type": "string",
"default": "root",
"description": (
'Folder scope (default "root"). Pass a specific folder '
'ID to scope into that folder, or "root" to reference '
"the library root. The read-only \"shared-with-me\" and "
'"following" folders live at the library root — pass '
"one of those ids to browse them. Copy any folder_id "
"verbatim from a browse/tree response, never construct "
"one. Combine with `recursive` to control breadth."
),
},
"recursive": {
"type": "boolean",
"default": False,
"description": (
"Whether to include documents from descendant folders. "
"When false (default), returns the direct contents of "
"folder_id along with its sub-folders — prefer this for "
"level-by-level exploration so you retain folder "
"hierarchy context. When true, flattens all descendant "
"documents into one list and omits sub-folders — use "
"only when a non-recursive browse of the target folder "
"returned no relevant results and you need to widen the "
"scope, or the user explicitly requests a flat listing."
),
},
"sort": {
"type": "string",
"enum": ["time", "relevance"],
"default": "time",
"description": (
'Sort order. "time" (default) sorts by upload date '
'(newest first); "relevance" orders documents by '
"semantic relevance to `query`. Relevance also works "
"inside the read-only shared folders — pass their "
"folder_id — but at the library root it ranks only "
"your own documents."
),
},
"query": {
"type": "string",
"description": (
"Search query for relevance ranking. Required when "
'sort="relevance"; must be omitted when sort="time".'
),
},
"offset": {
"type": "integer",
"minimum": 0,
"maximum": 9007199254740991,
"default": 0,
"description": (
"Zero-based pagination offset. Pass the value of "
"`next_offset` from the previous response to fetch the "
"next page."
),
},
"limit": {
"type": "number",
"minimum": 1,
"maximum": 50,
"default": 10,
"description": (
"Number of documents to return per page (1-50, "
"default 10)"
),
},
},
"required": [],
},
},
"get_document": {
"annotations": {"readOnlyHint": True, "openWorldHint": False},
"description": (
"Check a document's processing status and metadata. `status` is "
'one of "pending", "queued", "processing", "completed", or '
'"failed" — call this before `get_document_structure()` or '
"`get_page_content()` to confirm the document is ready."
),
"schema": {
"type": "object",
"properties": {
"doc_name": {
"type": "string",
"minLength": 1,
"description": _DOC_NAME_DESCRIPTION,
},
"folder_id": {
"anyOf": [{"type": "string"}, {"type": "null"}],
"description": _FOLDER_ID_DISAMBIGUATOR_DESCRIPTION,
},
"wait_for_completion": {
"type": "boolean",
"default": False,
"description": _WAIT_FOR_COMPLETION_DESCRIPTION,
},
},
"required": ["doc_name"],
},
},
"get_document_structure": {
"annotations": {"readOnlyHint": True, "openWorldHint": False},
"description": (
"Extract a document's hierarchical outline (headers, sections, "
f"page references). REQUIRED for documents over "
f"{STRUCTURE_FIRST_PAGE_THRESHOLD} pages — call this first to "
"locate relevant sections, then pass their page numbers to "
"`get_page_content()`. Use the `part` parameter to iterate large "
"outlines until `pagination.has_more` is false."
),
"schema": {
"type": "object",
"properties": {
"doc_name": {
"type": "string",
"minLength": 1,
"description": _DOC_NAME_DESCRIPTION,
},
"folder_id": {
"anyOf": [{"type": "string"}, {"type": "null"}],
"description": _FOLDER_ID_DISAMBIGUATOR_DESCRIPTION,
},
"part": {
"type": "integer",
"minimum": 1,
"maximum": 9007199254740991,
"default": 1,
"description": (
"Part number for pagination (1-based, default 1). For "
"large outlines, increment until the response's "
"`pagination.has_more` becomes false."
),
},
"wait_for_completion": {
"type": "boolean",
"default": False,
"description": _WAIT_FOR_COMPLETION_DESCRIPTION,
},
},
"required": ["doc_name"],
},
},
"get_page_content": {
"annotations": {"readOnlyHint": True, "openWorldHint": False},
"description": (
"Extract page content from a processed document. Use tight, "
"targeted page ranges — never the whole document at once. For "
f"documents over {STRUCTURE_FIRST_PAGE_THRESHOLD} pages, call "
"`get_document_structure()` first to pick relevant sections. "
"Embedded image paths in the response feed into "
"`get_document_image()`."
),
"schema": {
"type": "object",
"properties": {
"doc_name": {
"type": "string",
"minLength": 1,
"description": _DOC_NAME_DESCRIPTION,
},
"folder_id": {
"anyOf": [{"type": "string"}, {"type": "null"}],
"description": _FOLDER_ID_DISAMBIGUATOR_DESCRIPTION,
},
"pages": {
"type": "string",
"minLength": 1,
"pattern": r"^(\d+(-\d+)?)(,\s*\d+(-\d+)?)*$",
"description": (
'Page specification: "5", "3,7,10", "5-10", or '
'"1-3,7,9-12"'
),
},
"wait_for_completion": {
"type": "boolean",
"default": False,
"description": _WAIT_FOR_COMPLETION_DESCRIPTION,
},
},
"required": ["doc_name", "pages"],
},
},
"remove_document": {
"annotations": {"readOnlyHint": False, "destructiveHint": True,
"idempotentHint": True, "openWorldHint": False},
"description": (
"Permanently delete documents and all associated data. Only invoke "
"when the user explicitly names the documents AND confirms "
"deletion. Returns `results` — one entry per requested document: "
'`{ doc_name, status: "deleted" | "not_found" | "failed", '
"error? }`. Inspect each entry for per-document failures. This "
"action is irreversible."
),
"schema": {
"type": "object",
"properties": {
"doc_names": {
"type": "array",
"items": {"type": "string", "minLength": 1},
"minItems": 1,
"maxItems": 10,
"description": (
"Array of document names to delete. Each name must be "
"copied verbatim from the `name` field of a "
"browse_documents() or search_documents() response "
"(case-sensitive, include extension). Example: "
'["Q3 Report.pdf", "draft.pdf"]. Max 10 per call.'
),
},
"folder_id": {
"anyOf": [{"type": "string"}, {"type": "null"}],
"description": _FOLDER_ID_DISAMBIGUATOR_DESCRIPTION,
},
},
"required": ["doc_names"],
},
},
}
_READ_TOOLS = ("browse_documents", "get_document", "get_document_structure",
"get_page_content")
_MANAGEMENT_TOOLS = ("remove_document",)
# ── response envelopes ──
_ToolResult = tuple[dict, bool]
def _success(data: dict[str, Any], next_steps: dict[str, Any]) -> tuple[dict, bool]:
return {"success": True, **data, "next_steps": next_steps}, False
def _failure(error: str, details: Optional[dict[str, Any]],
next_steps: dict[str, Any], error_code: Optional[str] = None,
) -> tuple[dict, bool]:
payload: dict[str, Any] = {"error": error}
if error_code:
payload["errorCode"] = error_code
if details:
payload.update(details)
payload["next_steps"] = next_steps
return payload, True
def _dumps(payload: dict[str, Any]) -> str:
# Compact, matching _serialized_size — so the size budget measures
# what is actually emitted.
return json.dumps(payload, ensure_ascii=False)
# ── document listing / name resolution ──
def _all_documents(client) -> list[dict[str, Any]]:
"""Every document the client can list, newest first (both modes list
newest-first; paging preserves that order)."""
documents: list[dict[str, Any]] = []
offset = 0
while True:
page = client.list_documents(limit=100, offset=offset)
batch = page.get("documents") or []
documents.extend(batch)
# Advance by what actually arrived — stepping by the requested
# limit skips documents whenever a server caps its page size.
offset += len(batch)
total = page.get("total")
# An empty page is the reliable terminator; `total` (absent or
# None on some backends) only saves the final empty-page request.
if not batch or (isinstance(total, int) and offset >= total):
return documents
def _normalize_created_at(value: Any) -> str:
"""Emit the cloud tool format (ISO-8601 UTC with 'Z', millisecond
precision) from either mode's createdAt string."""
if not isinstance(value, str) or not value:
return ""
try:
from datetime import datetime, timezone
parsed = datetime.fromisoformat(value.replace("Z", "+00:00"))
if parsed.tzinfo is None:
parsed = parsed.replace(tzinfo=timezone.utc)
parsed = parsed.astimezone(timezone.utc)
return parsed.isoformat(timespec="milliseconds").replace("+00:00", "Z")
except ValueError:
return value
def _flat_metadata(value: Any) -> Optional[dict[str, Any]]:
"""User-facing string|number|boolean metadata fields only, or None."""
if not isinstance(value, dict):
return None
flat = {key: val for key, val in value.items()
if isinstance(val, (str, int, float, bool))}
return flat or None
def _scope_documents(documents: list[dict[str, Any]],
allowed_ids: Optional[frozenset]) -> list[dict[str, Any]]:
if allowed_ids is None:
return documents
return [doc for doc in documents if doc.get("id") in allowed_ids]
def _resolve_document(
client, doc_name: str,
documents: Optional[list[dict[str, Any]]] = None,
allowed_ids: Optional[frozenset] = None,
) -> "tuple[Optional[dict[str, Any]], Optional[_ToolResult]]":
"""Resolve doc_name to a list entry. Same-name duplicates resolve to the
newest match. Returns (entry, None) or (None, error_payload_pair)."""
if documents is None:
documents = _all_documents(client)
documents = _scope_documents(documents, allowed_ids)
matches = [doc for doc in documents if doc.get("name") == doc_name]
if matches:
return max(matches, key=lambda d: d.get("createdAt") or ""), None
names = [str(doc.get("name")) for doc in documents if doc.get("name")]
similar = difflib.get_close_matches(doc_name, names, n=_SIMILAR_NAMES_LIMIT,
cutoff=0.5)
message = (
"Document not found. Did you mean: "
+ ", ".join(f'"{name}"' for name in similar) + "?"
if similar else "Document not found or you do not have access to it"
)
return None, _failure(
message,
{"doc_name": doc_name, "similar_files": similar},
{
"summary": "The requested document does not exist or is not accessible",
"options": [
"Verify the document name is correct",
"Use browse_documents() to see your recent documents",
"Check if the document was deleted",
],
},
"NOT_FOUND",
)
def _refetch_entry(client, doc_id: str) -> Optional[dict[str, Any]]:
try:
return client.get_document(doc_id)
except PageIndexAPIError:
return None
def _await_completion(client, entry: dict[str, Any], wait: bool) -> dict[str, Any]:
"""Re-poll a processing document for up to 3 minutes when wait is set."""
doc_id = entry.get("id")
if not wait or not doc_id or entry.get("status") in ("completed", "failed"):
return entry
deadline = time.monotonic() + _TOOL_WAIT_TIMEOUT
current = entry
while time.monotonic() < deadline:
time.sleep(_TOOL_WAIT_INTERVAL)
refreshed = _refetch_entry(client, doc_id)
if refreshed is None:
return current
if refreshed.get("metadata") is None:
# Status refetches omit (or null out) custom metadata; keep the
# listing's copy.
refreshed["metadata"] = current.get("metadata")
current = {**current, **refreshed}
if current.get("status") in ("completed", "failed"):
return current
return current
def _not_ready_error(doc_name: str, status: Any, operation: str,
timed_out: bool) -> tuple[dict, bool]:
if status == "failed":
return _failure(
f"Document processing failed. Current status: {status}",
{"doc_name": doc_name},
{
"summary": "Document processing has failed",
"options": [
"Index the document again with "
"PageIndexClient.submit_document()",
"Use browse_documents() to work with other documents",
],
},
"INVALID_INPUT",
)
if timed_out:
return _failure(
f"Document is still processing. Current status: {status}",
{"doc_name": doc_name},
{
"summary": "Document processing timeout",
"options": [
"Try again later when processing is complete",
"Check status with get_document()",
],
},
"INVALID_INPUT",
)
return _failure(
f"Document is not ready for {operation}. Current status: {status}",
{"doc_name": doc_name},
{
"summary": "Document is still processing",
"options": [
"Wait for document processing to complete",
"Check status with browse_documents() or get_document()",
],
},
"INVALID_INPUT",
)
def _folder_unsupported(param: str) -> tuple[dict, bool]:
return _failure(
f"Folders are not supported in local mode yet — omit {param}.",
None,
{
"summary": "This local library does not have folders yet",
"options": ["Retry the call without a folder_id",
"Use browse_documents() to list the library root",
"Folders are available on PageIndex cloud (PageIndexCloudClient with an API key)"],
},
"INVALID_INPUT",
)
# ── page spec handling ──
def _parse_page_spec(
pages: str, doc_name: str,
) -> "tuple[Optional[list[int]], Optional[_ToolResult]]":
"""Expand '1-3,7' into a sorted, deduplicated page list, or an error."""
invalid = _failure(
"Invalid page specification format",
{"doc_name": doc_name},
{
"summary": "Failed to parse the pages parameter",
"options": [
'Use valid formats: "5", "3,7,10", "5-10", or "1-3,7,9-12"',
"Ensure page numbers are positive integers",
],
},
"INVALID_INPUT",
)
if not isinstance(pages, str) or not _PAGES_SPEC_RE.match(pages.strip()):
return None, invalid
too_many = _failure(
f"Too many pages requested (over {_MAX_REQUESTED_PAGES})",
{"doc_name": doc_name},
{
"summary": "The page specification spans too many pages",
"options": [
"Request a narrower page range",
"The response holds only a few pages per call - page through with several smaller requests",
],
},
"INVALID_INPUT",
)
expanded: set[int] = set()
for part in pages.split(","):
part = part.strip()
if "-" in part:
start, end = (int(x) for x in part.split("-", 1))
if start > end:
return None, invalid
else:
start = end = int(part)
# Bound each part arithmetically before materializing it: a spec like
# "1-1000000000" would otherwise expand to billions of integers
# inside the caller's process. The cap is on distinct pages, so
# overlapping parts (a parent section plus its children) don't
# double-count.
if end - start + 1 > _MAX_REQUESTED_PAGES:
return None, too_many
expanded.update(range(start, end + 1))
if len(expanded) > _MAX_REQUESTED_PAGES:
return None, too_many
if any(page < 1 for page in expanded):
return None, _failure(
"Invalid page numbers. Page numbers must be positive integers",
{"doc_name": doc_name},
{
"summary": "Invalid page numbers provided",
"options": [
"Page numbers must be positive integers (>= 1)",
"Check the page specification format",
],
},
"INVALID_INPUT",
)
return sorted(expanded), None
def _format_page_spec(pages: list[int]) -> str:
"""Compress [1,2,3,5] into '1-3,5'."""
if not pages:
return ""
ordered = sorted(set(pages))
ranges = []
start = prev = ordered[0]
for page in ordered[1:]:
if page == prev + 1:
prev = page
continue
ranges.append(f"{start}" if start == prev else f"{start}-{prev}")
start = prev = page
ranges.append(f"{start}" if start == prev else f"{start}-{prev}")
return ",".join(ranges)
# ── structure formatting / splitting ──
_STRUCTURE_KEY_ORDER = ("title", "node_id", "start_index", "end_index",
"page_index", "prefix_summary", "summary", "nodes")
def _format_structure(node: Any) -> Any:
"""Drop node text and normalize key order, recursively."""
if isinstance(node, list):
return [_format_structure(item) for item in node]
if isinstance(node, dict):
stripped = {key: value for key, value in node.items() if key != "text"}
if "nodes" in stripped:
stripped["nodes"] = _format_structure(stripped["nodes"])
ordered = {key: stripped[key] for key in _STRUCTURE_KEY_ORDER
if key in stripped}
ordered.update({key: value for key, value in stripped.items()
if key not in ordered})
return ordered
return node
def _serialized_size(value: Any) -> int:
return len(json.dumps(value, ensure_ascii=False))
def _split_structure(structure: Any, budget: int) -> list[Any]:
"""Split a formatted structure into chunks of at most ~budget serialized
chars. The paginated response shape matches the cloud tool (its chunk
type admits node-or-list); chunk boundaries are implementation-defined.
An unsplit structure keeps its natural shape; once split, every chunk
is a list of nodes — the `structure` field must not change JSON type
between parts of one paginated response."""
if _serialized_size(structure) <= budget:
return [structure]
nodes = structure if isinstance(structure, list) else [structure]
chunks: list[Any] = []
group: list[Any] = []
group_size = 0
for node in nodes:
size = _serialized_size(node)
if size > budget:
if group:
chunks.append(group)
group, group_size = [], 0
chunks.extend([part]
for part in _split_oversized_node(node, budget))
continue
if group and group_size + size > budget:
chunks.append(group)
group, group_size = [], 0
group.append(node)
group_size += size
if group:
chunks.append(group)
return chunks or [structure]
def _split_oversized_node(node: Any, budget: int) -> list[Any]:
children = node.get("nodes") if isinstance(node, dict) else None
if not children:
return [node]
shell = {key: value for key, value in node.items() if key != "nodes"}
shell_size = _serialized_size(shell)
child_budget = max(budget - shell_size, budget // 2)
parts = []
for chunk in _split_structure(children, child_budget):
# A recursive result is either the unsplit children (natural shape)
# or always-list chunks; normalize for the shell's "nodes".
parts.append({**shell,
"nodes": chunk if isinstance(chunk, list) else [chunk]})
return parts
# ── tool implementations (client-backed; mode-blind) ──
def _browse_documents(client, folder_id: str = "root", recursive: bool = False,
sort: str = "time", query: Optional[str] = None,
offset: int = 0, limit: int = 10,
_allowed_ids: Optional[frozenset] = None) -> tuple[dict, bool]:
if folder_id != "root":
return _folder_unsupported("folder_id")
if sort not in ("time", "relevance"):
return _failure(
'Invalid sort mode — only the default "time" sort is available '
"in local mode.", None,
{"summary": "Invalid sort mode",
"options": ['Use sort="time" (newest first) or omit sort',
"Semantic ranking is available on PageIndex cloud (PageIndexCloudClient with an API key)"]},
"INVALID_INPUT",
)
if sort == "relevance" or query:
# Semantic ranking is a cloud capability; like folders, it is not
# imitated here.
return _failure(
"Relevance ranking is not supported in local mode yet — use "
"the default time sort.", None,
{"summary": "This local library does not have semantic ranking yet",
"options": ["Retry without sort/query and match the returned names and descriptions against the intent yourself",
"Page through the full library with `offset: next_offset`",
"Semantic ranking is available on PageIndex cloud (PageIndexCloudClient with an API key)"]},
"INVALID_INPUT",
)
try:
offset = max(int(offset), 0)
limit = min(max(int(limit), 1), 50)
except (TypeError, ValueError):
return _failure("offset and limit must be numbers", None,
{"summary": "Invalid pagination parameters",
"options": ["Pass integer offset and limit values"]},
"INVALID_INPUT")
if _allowed_ids is None:
listing = client.list_documents(limit=limit, offset=offset)
window = listing.get("documents") or []
total = listing.get("total")
else:
scoped = _scope_documents(_all_documents(client), _allowed_ids)
window, total = scoped[offset:offset + limit], len(scoped)
# Advance by what actually arrived — a server may cap its page size —
# and treat an absent/None total like _all_documents does: a full
# window means there may be more.
window_end = offset + len(window)
has_more = bool(window) and (window_end < total if isinstance(total, int)
else len(window) == limit)
next_offset = window_end if has_more else None
page_has_processing = False
page_has_failed = False
items = []
for doc in window:
status = doc.get("status") or "unknown"
if status == "failed":
page_has_failed = True
elif status != "completed":
page_has_processing = True
item = {
"name": doc.get("name") or "Unknown Document",
"description": doc.get("description") or "No description provided",
"status": status,
"created_at": _normalize_created_at(doc.get("createdAt")),
}
metadata = _flat_metadata(doc.get("metadata"))
if metadata is not None:
item["metadata"] = metadata
items.append(item)
data: dict[str, Any] = {
"documents": items,
"sort": sort,
"next_offset": next_offset,
"has_more": has_more,
}
if not recursive:
data["folders"] = []
if not items and offset == 0:
next_steps = {
"summary": "Nothing to show",
"options": ["Nothing here. Index documents with "
"PageIndexClient.submit_document() to get started."],
"auto_retry": "Index a document with "
"PageIndexClient.submit_document() to get started",
}
return _success(data, next_steps)
options = []
if items:
options.append("Use get_document() with a document name to view details")
options.append(
"Results returned ≠ correct results. Verify these documents match "
"the user's actual intent (topic, time period, document type) "
"before proceeding."
+ (" If they do not match, page through the rest of the library."
if has_more else "")
+ " Do NOT use general knowledge as a substitute."
)
if page_has_processing:
options.append("Some documents on this page are still processing. "
"Use get_document() to check individual status.")
if page_has_failed:
options.append("Some documents on this page failed processing. "
"Use get_document() to see error details.")
if has_more:
options.append("Use browse_documents() with `offset: next_offset` to "
"load more documents")
summary = (f"Showing {len(items)} document(s)"
+ (" (more available)" if has_more else "")
if items else "Nothing to show")
return _success(data, {"summary": summary, "options": options})
def _get_document(client, doc_name: str, folder_id: Optional[str] = None,
wait_for_completion: bool = False,
_allowed_ids: Optional[frozenset] = None) -> tuple[dict, bool]:
if folder_id not in (None, "root"):
return _folder_unsupported("folder_id")
entry, error = _resolve_document(client, doc_name, allowed_ids=_allowed_ids)
if error is not None:
return error
assert entry is not None
entry = _await_completion(client, entry, wait_for_completion)
status = entry.get("status") or "unknown"
is_processing = status not in ("completed", "failed")
is_ready = status == "completed"
page_num = entry.get("pageNum") or 0
name = entry.get("name") or "Unknown Document"
suggestions: list[str] = []
if is_processing:
suggestions.append("Document is still processing. Processing status "
"can be checked later.")
elif is_ready:
suggestions.append("Document is ready for analysis.")
if page_num > 0:
if page_num <= 5:
suggestions.extend([
f"This is a short document with {page_num} pages.",
f'First explore structure: get_document_structure(doc_name: "{name}")',
f'Then extract all content: get_page_content(doc_name: "{name}", pages: "1-{page_num}")',
])
elif page_num <= STRUCTURE_FIRST_PAGE_THRESHOLD:
suggestions.extend([
f"This document has {page_num} pages.",
f'First explore structure: get_document_structure(doc_name: "{name}")',
f'Then extract key pages: get_page_content(doc_name: "{name}", pages: "1,5,10")',
])
else:
suggestions.extend([
f"This is a large document with {page_num} pages.",
f'First explore structure: get_document_structure(doc_name: "{name}")',
f'Then target specific sections: get_page_content(doc_name: "{name}", pages: "1-3")',
])
else:
suggestions.append("Document processing failed. Index the document "
"again with PageIndexClient.submit_document().")
data: dict[str, Any] = {
"name": name,
"description": entry.get("description") or "No description provided",
"status": status,
"created_at": _normalize_created_at(entry.get("createdAt")),
"page_count": page_num or None,
"folder_id": entry.get("folderId"),
}
metadata = _flat_metadata(entry.get("metadata"))
if metadata is not None:
data["metadata"] = metadata
return _success(data, {
"summary": ("Document is ready for analysis and querying." if is_ready
else "Document is still being processed." if is_processing
else "Document processing has failed."),
"options": suggestions,
**({"auto_retry": "Document processing status can be monitored periodically"}
if is_processing else {}),
})
def _get_document_structure(client, doc_name: str,
folder_id: Optional[str] = None, part: int = 1,
wait_for_completion: bool = False,
_allowed_ids: Optional[frozenset] = None) -> tuple[dict, bool]:
if folder_id not in (None, "root"):
return _folder_unsupported("folder_id")
entry, error = _resolve_document(client, doc_name, allowed_ids=_allowed_ids)
if error is not None:
return error
assert entry is not None
waited = wait_for_completion and entry.get("status") not in ("completed", "failed")
entry = _await_completion(client, entry, wait_for_completion)
if entry.get("status") != "completed":
return _not_ready_error(doc_name, entry.get("status"),
"structure retrieval",
waited and entry.get("status") != "failed")
try:
# Prefer the raw stored tree: its nodes carry start_index/end_index
# like the cloud structure tool, where client.get_tree() drops
# end_index and renames fields.
raw_tree = getattr(getattr(client, "_api", None), "raw_tree", None)
tree = raw_tree(entry["id"]) if raw_tree is not None else None
if tree is None:
tree = client.get_tree(entry["id"], node_summary=True).get("result")
except PageIndexAPIError as exc:
return _failure(
f"Failed to retrieve document structure: {exc}",
{"doc_name": doc_name},
{
"summary": "Failed to retrieve document structure due to an error",
"options": [
"The document may not exist or is not accessible",
"Check if the document name is correct",
"Try again in a few moments",
],
},
"INTERNAL_ERROR",
)
if tree is None:
return _failure(
"Structure not available for this document",
{"doc_name": doc_name},
{
"summary": "Structure not available for this document",
"options": [
"The document may not have been processed correctly or structure extraction may have failed",
"Try processing the document again if possible",
],
},
"INTERNAL_ERROR",
)
formatted = _format_structure(tree)
chunks = _split_structure(formatted, _CHAR_BUDGET)
total_parts = max(1, len(chunks))
try:
requested_part = int(part)
except (TypeError, ValueError):
requested_part = 1
current = min(max(requested_part, 1), total_parts)
if total_parts == 1:
return _success(
{"doc_name": doc_name, "structure": chunks[0]},
{
"summary": "Document structure retrieved successfully.",
"options": [
"Use get_page_content() to extract specific content from pages",
],
},
)
next_steps = (
{
"summary": f"Showing part {current} of {total_parts}.",
"options": [
f"Request next part with part: {current + 1}",
f"Jump to last part with part: {total_parts}",
"Proceed to get_page_content() for specific sections",
],
}
if current < total_parts else
{
"summary": "All parts retrieved for current pagination.",
"options": [
"Use get_page_content() to extract specific content from pages",
],
}
)
return _success(
{
"doc_name": doc_name,
"total_parts": total_parts,
"structure": chunks[current - 1],
"pagination": {
"part": current,
"total_parts": total_parts,
"has_more": current < total_parts,
},
},
next_steps,
)
def _get_page_content(client, doc_name: str, pages: str,
folder_id: Optional[str] = None,
wait_for_completion: bool = False,
_allowed_ids: Optional[frozenset] = None) -> tuple[dict, bool]:
if folder_id not in (None, "root"):
return _folder_unsupported("folder_id")
entry, error = _resolve_document(client, doc_name, allowed_ids=_allowed_ids)
if error is not None:
return error
assert entry is not None
waited = wait_for_completion and entry.get("status") not in ("completed", "failed")
entry = _await_completion(client, entry, wait_for_completion)
if entry.get("status") != "completed":
return _not_ready_error(doc_name, entry.get("status"),
"page content retrieval",
waited and entry.get("status") != "failed")
requested, error = _parse_page_spec(pages, doc_name)
if error is not None:
return error
assert requested is not None
try:
page_data = client.get_ocr(entry["id"], format="page").get("result") or []
except PageIndexAPIError as exc:
return _failure(
f"Failed to retrieve page content: {exc}",
{"doc_name": doc_name},
{
"summary": "Unable to retrieve page content due to a service issue.",
"options": [
"Verify the document name is correct using browse_documents()",
"Check if the document processing is complete with get_document()",
"Ensure the requested page numbers are valid",
],
"auto_retry": "This may be a temporary issue - you can try "
"the request again",
},
"INTERNAL_ERROR",
)
by_index = {item["page_index"]: item for item in page_data
if isinstance(item, dict)
and isinstance(item.get("page_index"), int)}
max_page = max(by_index, default=0)
out_of_range = [page for page in requested if page > max_page]
valid_pages = [page for page in requested if page <= max_page]
if out_of_range and not valid_pages:
return _failure(
f"All requested pages are out of range. Document has {max_page} "
f"pages, but you requested pages: {_format_page_spec(out_of_range)}",
{
"doc_name": doc_name,
"max_pages": max_page,
"requested_pages": _format_page_spec(out_of_range),
},
{
"summary": "All requested pages are out of range for this document",
"options": [
f"Request pages between 1 and {max_page}",
"Use get_document() to check document page count",
],
},
"INVALID_INPUT",
)
content = []
included: list[int] = []
remaining: list[int] = []
budget = _CHAR_BUDGET
for page in valid_pages:
item = by_index.get(page)
markdown = item.get("markdown") if item else None
text = (markdown if isinstance(markdown, str)
else f"Page {page} content not available")
if not included or budget - len(text) >= 0:
content.append({"page": page, "text": text})
included.append(page)
budget -= len(text)
else:
remaining.append(page)
options = [
"Use get_document_structure() to understand document organization",
"Request additional pages as needed",
]
if remaining:
options.insert(0, f"For remaining pages, request: {_format_page_spec(remaining)}")
if out_of_range:
options.insert(0, f"Document has {max_page} pages total - request "
f"pages 1-{max_page}")
# Additive, not either/or: a call can both truncate for size and have
# out-of-range pages — hiding either would misreport what was returned.
if remaining or out_of_range:
parts = [f"Retrieved {len(included)} of {len(requested)} "
"requested pages."]
if remaining:
parts.append(f"Pages {_format_page_spec(remaining)} were "
"omitted due to response size limits.")
if out_of_range:
parts.append(f"Pages {_format_page_spec(out_of_range)} "
"were out of range.")
summary = " ".join(parts)
else:
summary = (f"Successfully retrieved content for {len(content)} "
f"page{'' if len(content) == 1 else 's'}.")
return _success(
{
"doc_name": doc_name,
"total_pages": max_page,
"requested_pages": _format_page_spec(requested),
"returned_pages": _format_page_spec(included),
"content": content,
},
{"summary": summary, "options": options},
)
def _remove_document(client, doc_names: list[str],
folder_id: Optional[str] = None,
_allowed_ids: Optional[frozenset] = None) -> tuple[dict, bool]:
if folder_id not in (None, "root"):
return _folder_unsupported("folder_id")
if not isinstance(doc_names, list) or not doc_names:
return _failure("At least one document name is required", None,
{"summary": "No document names provided",
"options": ["Pass doc_names as a non-empty array"]},
"INVALID_INPUT")
# Validate every element before deleting anything: a rejection envelope
# must mean nothing was destroyed.
if not all(isinstance(name, str) and name.strip() for name in doc_names):
return _failure(
"doc_names must be an array of non-empty document name strings",
None,
{"summary": "Invalid document names",
"options": ["Copy each name verbatim from a browse_documents() "
"response"]},
"INVALID_INPUT")
# A repeated name is one deletion, not a second "failed" row.
doc_names = list(dict.fromkeys(doc_names))
if len(doc_names) > 10:
return _failure("Maximum 10 documents can be deleted at once", None,
{"summary": "Too many documents in one call",
"options": ["Delete at most 10 documents per call"]},
"INVALID_INPUT")
documents = _all_documents(client)
results = []
for doc_name in doc_names:
entry, error = _resolve_document(client, doc_name, documents=documents,
allowed_ids=_allowed_ids)
if error is not None or entry is None:
results.append({"doc_name": doc_name, "status": "not_found"})
continue
try:
client.delete_document(entry["id"])
results.append({"doc_name": doc_name, "status": "deleted"})
except Exception as exc:
# Any escape here (OSError, transport errors) would discard the
# entries for documents already irreversibly deleted.
results.append({"doc_name": doc_name, "status": "failed",
"error": str(exc)})
deleted = sum(1 for item in results if item["status"] == "deleted")
return _success(
{"results": results},
{
"summary": f"Deleted {deleted} of {len(doc_names)} document(s).",
"options": ["Use browse_documents() to review the remaining library"],
},
)
_IMPLEMENTATIONS: dict[str, Callable[..., tuple[dict, bool]]] = {
"browse_documents": _browse_documents,
"get_document": _get_document,
"get_document_structure": _get_document_structure,
"get_page_content": _get_page_content,
"remove_document": _remove_document,
}
def tool_names(include_management: bool = False) -> tuple[str, ...]:
return _READ_TOOLS + (_MANAGEMENT_TOOLS if include_management else ())
def _coerce_bool_args(name: str, kwargs: dict[str, Any]) -> None:
"""Models routinely send booleans as JSON strings ("false"); the bare
truthiness tests downstream would read those as True."""
properties = TOOL_CONTRACT.get(name, {}).get("schema", {}).get(
"properties", {})
for key, spec in properties.items():
value = kwargs.get(key)
if spec.get("type") == "boolean" and isinstance(value, str):
kwargs[key] = value.strip().lower() not in ("false", "no", "0", "")
def call_tool(client, name: str, arguments: dict[str, Any],
doc_ids=None) -> tuple[str, bool]:
"""Run one contract tool; returns (envelope_json, is_error). Never raises
for tool-level failures — unexpected exceptions become error envelopes.
``doc_ids`` restricts every document lookup to that allowlist (the local
chat surfaces' doc_id scope)."""
implementation = _IMPLEMENTATIONS.get(name)
if implementation is None:
payload, _ = _failure(
f"Unknown tool: {name}",
{"tool_name": name, "available_tools": list(_IMPLEMENTATIONS)},
{"summary": "Tool not found",
"options": [f"Available tools: {', '.join(_IMPLEMENTATIONS)}"]},
"INVALID_INPUT",
)
return _dumps(payload), True
if arguments is not None and not isinstance(arguments, dict):
payload, is_error = _failure(
f"Invalid arguments for {name}: expected a JSON object, got "
f"{type(arguments).__name__}", None,
{"summary": "Invalid tool arguments",
"options": [f"Pass {name}() arguments as a JSON object of its "
"parameters"]},
"INVALID_INPUT",
)
return _dumps(payload), is_error
# Underscore-prefixed keys are the SDK's private channel (the scope
# below), never model arguments. None ≡ omitted (the contract's
# "omit if ..." semantics, same as the cloud bridge invoker).
kwargs = {key: value for key, value in (arguments or {}).items()
if not key.startswith("_") and value is not None}
_coerce_bool_args(name, kwargs)
try:
if doc_ids is not None:
ids = [doc_ids] if isinstance(doc_ids, str) else doc_ids
kwargs["_allowed_ids"] = frozenset(str(one_id) for one_id in ids)
bound = inspect.signature(implementation).bind(client, **kwargs)
except TypeError as exc:
payload, is_error = _failure(
f"Invalid arguments for {name}: {exc}", None,
{"summary": "Invalid tool arguments",
"options": [f"Check the {name}() parameter names and types"]},
"INVALID_INPUT",
)
return _dumps(payload), is_error
try:
payload, is_error = implementation(*bound.args, **bound.kwargs)
except Exception as exc: # tool calls must never raise into the agent loop
payload, is_error = _failure(
f"{name} failed: {exc}", None,
{"summary": "Unexpected error while running the tool",
"options": ["Try the request again"],
"auto_retry": "This is likely a temporary issue - you can try "
"the request again"},
"INTERNAL_ERROR",
)
return _dumps(payload), is_error
# ── plain-function materialization (the `client.agent_tools()` surface) ──
def _tool_docstring(description: str, properties: dict[str, Any]) -> str:
lines = [description, "", "Args:"]
for param, spec in properties.items():
lines.append(f" {param}: {spec.get('description', '')}")
return "\n".join(lines)
# Local guidance layer: schema STRUCTURE stays byte-identical to the cloud
# contract minus the hidden cloud-only parameters, and description strings
# adapt to the local surface the same way AGENT_INSTRUCTIONS does — guidance
# must not teach capabilities (folders, semantic ranking) or tools
# (search_documents, get_document_image) that do not exist here. Guard
# tests pin structure (contract-minus-hidden equality), tool references
# (the dead-reference test), and capability phrases (the per-docstring
# phrase test) — a contract refresh that reintroduces a cloud-only
# reference fails loudly.
#: Cloud-only parameters hidden from the local surface — strict-schema
#: frameworks make the dead-end calls inexpressible, and lenient framework
#: argument models drop them before the call (degrading to the bare call).
#: The call_tool path still answers folder_id/sort/query with the guided
#: error envelope; recursive is simply accepted (flattening a folderless
#: library is the identity). Plain functions reject unknown parameters at
#: the Python call boundary.
_LOCAL_HIDDEN_PARAMS: dict[str, tuple[str, ...]] = {
"browse_documents": ("folder_id", "recursive", "sort", "query"),
"get_document": ("folder_id",),
"get_document_structure": ("folder_id",),
"get_page_content": ("folder_id",),
"remove_document": ("folder_id",),
}
_LOCAL_DOC_NAME_DESCRIPTION = (
'Copy the `name` field verbatim from a browse_documents() response '
'(case-sensitive, include extension). Example: "Q3 Report.pdf". '
"Document names are unique in a local library."
)
_LOCAL_DESCRIPTIONS: dict[str, str] = {
"browse_documents": (
"Primary document retrieval tool — first choice for any "
"document-related question. Lists your documents newest first with "
"names and descriptions; match them against the user's intent and "
"page through with `offset: next_offset` (limit up to 50) while "
"`has_more` is true. "
'Folder browsing and semantic ranking (sort="relevance") are not '
"supported in local mode yet — they work on PageIndex cloud."
),
# The image sentence points at a tool that is not registered locally.
"get_page_content": TOOL_CONTRACT["get_page_content"]["description"]
.replace(" Embedded image paths in the response feed into "
"`get_document_image()`.", ""),
}
_LOCAL_PARAM_DESCRIPTIONS: dict[tuple[str, str], str] = {
("get_document", "doc_name"): _LOCAL_DOC_NAME_DESCRIPTION,
("get_document_structure", "doc_name"): _LOCAL_DOC_NAME_DESCRIPTION,
("get_page_content", "doc_name"): _LOCAL_DOC_NAME_DESCRIPTION,
("remove_document", "doc_names"): (
"Array of document names to delete. Each name must be copied "
"verbatim from the `name` field of a browse_documents() response "
'(case-sensitive, include extension). Example: ["Q3 Report.pdf", '
'"draft.pdf"]. Max 10 per call.'
),
}
def _local_description(name: str) -> str:
return _LOCAL_DESCRIPTIONS.get(name) or TOOL_CONTRACT[name]["description"]
def _local_schema(name: str) -> dict[str, Any]:
schema = copy.deepcopy(TOOL_CONTRACT[name]["schema"])
for param in _LOCAL_HIDDEN_PARAMS.get(name, ()):
schema["properties"].pop(param, None)
for (tool_name, param), text in _LOCAL_PARAM_DESCRIPTIONS.items():
if tool_name == name and param in schema["properties"]:
schema["properties"][param]["description"] = text
return schema
def _docstring(name: str) -> str:
return _tool_docstring(_local_description(name),
_local_schema(name)["properties"])
_SCHEMA_TYPE_MAP = {"string": str, "integer": int, "number": float,
"boolean": bool, "array": list, "object": dict}
def _annotation_for(spec: dict) -> Any:
schema_type = spec.get("type")
if schema_type is None and isinstance(spec.get("anyOf"), list):
# Nullable unions arrive as anyOf: [{type: string}, {type: null}].
options = [option for option in spec["anyOf"]
if isinstance(option, dict) and option.get("type")]
schema_type = [option["type"] for option in options]
# `items` lives on the array option, not the union shell.
spec = next((option for option in options
if option["type"] == "array"), spec)
nullable = False
if isinstance(schema_type, list):
nullable = "null" in schema_type
bases = [t for t in schema_type if t != "null"]
schema_type = bases[0] if bases else None
base = _SCHEMA_TYPE_MAP.get(schema_type or "", Any)
if base is list:
# Strict function calling rejects arrays whose item type was lost
# in the annotation round-trip; parameterize when it is known.
item_type = (spec["items"].get("type")
if isinstance(spec.get("items"), dict) else None)
element = (_SCHEMA_TYPE_MAP.get(item_type)
if isinstance(item_type, str) else None)
if element is not None:
base = list[element]
return Optional[base] if nullable else base
def _bridge_invoker(bridge, name: str) -> "Callable[[dict], tuple[str, bool]]":
"""One cloud tool call proxied over MCP: None-valued arguments are
dropped (None ≡ omitted, matching the contract's "omit if ..."
semantics) and failures are contained in the error envelope. Returns
(envelope_text, is_error), like call_tool."""
def _invoke(arguments: dict[str, Any]) -> tuple[str, bool]:
try:
arguments = {key: value for key, value in arguments.items()
if value is not None}
return bridge.call_tool(name, arguments)
except Exception as exc:
payload, _ = _failure(
f"{name} failed: {exc}", None,
{"summary": "Unexpected error while running the tool",
"options": ["Try the request again"],
"auto_retry": "This is likely a temporary issue - you can "
"try the request again"},
"INTERNAL_ERROR",
)
return _dumps(payload), True
return _invoke
def _make_bridge_function(bridge, meta: dict) -> Callable[..., str]:
"""One plain function for a cloud tool: real signature and docstring from
the server's schema, invocation proxied over MCP, errors contained."""
import keyword
name = str(meta.get("name") or "")
schema = meta.get("inputSchema") or {}
properties: dict[str, Any] = schema.get("properties") or {}
required = set(schema.get("required") or [])
_invoke = _bridge_invoker(bridge, name)
params_usable = all(param.isidentifier() and not keyword.iskeyword(param)
and param != "_invoke"
for param in properties)
if not params_usable:
def proxy(**kwargs: Any) -> str:
return _invoke(kwargs)[0]
else:
ordered = ([p for p in properties if p in required]
+ [p for p in properties if p not in required])
rendered = ", ".join(
p if p in required else f"{p}={properties[p].get('default')!r}"
for p in ordered
)
args_literal = "{" + ", ".join(f"'{p}': {p}" for p in ordered) + "}"
namespace: dict[str, Any] = {"_invoke": _invoke}
exec(f"def _synthesized({rendered}):\n"
f" return _invoke({args_literal})[0]", namespace)
proxy = namespace["_synthesized"]
annotations: dict[str, Any] = {}
for p in ordered:
annotation = _annotation_for(properties[p])
if p not in required and "default" not in properties[p]:
# Absent-but-non-nullable params must admit None, or strict
# schemas force the model to always send a value.
annotation = Optional[annotation]
annotations[p] = annotation
annotations["return"] = str
proxy.__annotations__ = annotations
proxy.__name__ = proxy.__qualname__ = name or "tool"
proxy.__doc__ = _tool_docstring(meta.get("description") or "", properties)
return proxy
_BRIDGES: "weakref.WeakKeyDictionary" = weakref.WeakKeyDictionary()
_BRIDGES_LOCK = threading.Lock()
def _cloud_bridge(client):
"""One bridge per client: tool discovery and instructions share a single
MCP session. Weak-keyed off the instance so clients stay picklable; the
lock closes the check-then-set race under concurrent first calls."""
with _BRIDGES_LOCK:
bridge = _BRIDGES.get(client)
if bridge is None:
from .mcp_bridge import McpBridge
bridge = McpBridge(
f"{client.BASE_URL}/mcp",
{"Authorization": f"Bearer {client.api_key}"},
)
_BRIDGES[client] = bridge
return bridge
def _read_only_tools(tools_meta: list[dict]) -> list[dict]:
"""The management gate for consumers without a framework permission
layer: only tools the server marks read-only, guarded against a server
annotation regression silently disabling every tool."""
filtered = [meta for meta in tools_meta
if (meta.get("annotations") or {}).get("readOnlyHint") is True]
if tools_meta and not filtered:
raise PageIndexAPIError(
"The MCP server returned tools but none are annotated "
"read-only — a server annotation regression would otherwise "
"silently disable every tool. Pass include_management=True "
"to expose the unfiltered list."
)
return filtered
def _build_cloud_agent_tools(client, include_management: bool) -> list[Callable[..., str]]:
bridge = _cloud_bridge(client)
tools_meta = bridge.list_tools()
if not include_management:
tools_meta = _read_only_tools(tools_meta)
return [_make_bridge_function(bridge, meta) for meta in tools_meta]
def _require_local_scope(client, doc_ids) -> None:
"""The allowlist is enforced in-process; cloud lookups run server-side,
so accepting doc_ids there would be advisory-only — refuse loudly."""
if doc_ids is not None and getattr(client, "api_key", None):
raise PageIndexAPIError(
"doc_ids scoping applies to local tools only — cloud calls "
"are scoped server-side."
)
def _tool_specs(client, include_management: bool = False, doc_ids=None,
) -> "list[tuple[str, str, dict, Callable[[dict], tuple[str, bool]]]]":
"""(name, description, schema, invoke) per tool, for adapters that take
the wire schema verbatim. ``invoke`` returns (envelope_text, is_error).
Schemas are copies (frameworks keep the dict by reference). ``doc_ids``
is the local chat scope; cloud scoping is server-side."""
_require_local_scope(client, doc_ids)
if getattr(client, "api_key", None):
bridge = _cloud_bridge(client)
tools_meta = bridge.list_tools()
if not include_management:
tools_meta = _read_only_tools(tools_meta)
return [(str(meta.get("name") or "tool"),
meta.get("description") or "",
copy.deepcopy(meta.get("inputSchema"))
or {"type": "object", "properties": {}},
_bridge_invoker(bridge, str(meta.get("name") or "tool")))
for meta in tools_meta]
def local_invoke(name: str) -> "Callable[[dict], tuple[str, bool]]":
def invoke(arguments: dict) -> tuple[str, bool]:
return call_tool(client, name, arguments, doc_ids=doc_ids)
return invoke
return [(name, _local_description(name), _local_schema(name),
local_invoke(name))
for name in tool_names(include_management)]
def build_agent_tools(client, include_management: bool = False) -> list[Callable[..., str]]:
"""Plain synchronous functions bound to `client`.
Cloud: one function per tool of the live cloud MCP tool set, signatures
synthesized from the server's schemas, calls proxied over MCP. Local:
the built-in contract tools over the local store. Every function returns
the JSON envelope as a string and never raises for arguments its
signature accepts (cloud-only parameters are absent from the local
signatures; the call_tool path answers them with the guided envelope).
"""
if getattr(client, "api_key", None):
return _build_cloud_agent_tools(client, include_management)
def browse_documents(offset: int = 0, limit: int = 10) -> str:
return call_tool(client, "browse_documents", {
"offset": offset, "limit": limit,
})[0]
def get_document(doc_name: str, wait_for_completion: bool = False) -> str:
return call_tool(client, "get_document", {
"doc_name": doc_name,
"wait_for_completion": wait_for_completion,
})[0]
def get_document_structure(doc_name: str, part: int = 1,
wait_for_completion: bool = False) -> str:
return call_tool(client, "get_document_structure", {
"doc_name": doc_name, "part": part,
"wait_for_completion": wait_for_completion,
})[0]
def get_page_content(doc_name: str, pages: str,
wait_for_completion: bool = False) -> str:
return call_tool(client, "get_page_content", {
"doc_name": doc_name, "pages": pages,
"wait_for_completion": wait_for_completion,
})[0]
def remove_document(doc_names: list[str]) -> str:
return call_tool(client, "remove_document", {
"doc_names": doc_names,
})[0]
functions = {
"browse_documents": browse_documents,
"get_document": get_document,
"get_document_structure": get_document_structure,
"get_page_content": get_page_content,
"remove_document": remove_document,
}
tools = []
for name in tool_names(include_management):
function = functions[name]
function.__doc__ = _docstring(name)
tools.append(function)
return tools
# ── agent instructions ──
# Local subset of the cloud MCP server's initialize instructions (its
# no-folders variant), trimmed to what exists here: the search_documents
# escalation steps, get_document_image, and the shared read-only-folders
# block are removed, and the sort="relevance" guidance is replaced with
# name/description matching (semantic ranking is cloud-side). Cloud
# clients receive the server's live instructions instead — see
# _base_instructions().
_INSTRUCTIONS_HEADER = (
"PageIndex by Vectify AI is a document platform for uploading and "
"managing long PDFs (research papers, financial reports, legal docs, "
"textbooks, etc.)."
)
_READING_WORKFLOW = f"""\
READING WORKFLOW:
- For documents over {STRUCTURE_FIRST_PAGE_THRESHOLD} pages: call get_document_structure() first to locate relevant sections, then get_page_content() with targeted page ranges.
- For small documents ({STRUCTURE_FIRST_PAGE_THRESHOLD} pages or fewer): call get_page_content() directly."""
_TOOL_USAGE_RULES = """\
TOOL USAGE RULES:
- Invoke a tool only when all required parameters are present or clearly inferable. Never invent placeholder values.
- If a tool returns an error, present the provided next_steps/options to the user instead of retrying blindly."""
_DISCOVERY = """\
DOCUMENT DISCOVERY:
- browse_documents() — DEFAULT discovery tool, first choice for any document-related question. It lists your documents newest first with names and descriptions; match them against the user's intent, and page through with `offset: next_offset` while has_more is true."""
_DECISION = """\
DECISION:
- "What do I have / list / recent" → browse_documents()
- ANY question that needs a document to answer (including "find THE paper about Y") → browse_documents(), then pick the documents whose name/description matches the question"""
_AFTER_DISCOVERY = """\
- Skip discovery ONLY for questions with NO possible document connection (e.g., "capital of France").
- After discovery: 1 match or 1 clearly best match → proceed to read and answer without asking. Multiple equally relevant → ask user to pick.
- Results returned ≠ correct results. If the returned documents do not clearly match the user's intent (e.g., wrong topic, wrong time period, wrong document type), treat it the same as "not found" and continue the PERSISTENCE protocol below."""
_PERSISTENCE = """\
PERSISTENCE (before concluding the target document is not in the library):
This protocol applies both when results are empty AND when results are returned but none match the user's intent. Do NOT give up after a single discovery attempt. Follow these steps in order:
1. browse_documents() and compare every returned name/description against the user's intent
2. Page through the ENTIRE library with `limit: 50` and `offset: next_offset` until has_more is false — MANDATORY, must be completed before concluding "not found"
3. Re-scan for loose matches: synonyms, abbreviations, and partial titles in names/descriptions can identify the target
Only after ALL three steps have been tried may you conclude the document is not in the library. Do NOT fall back to general knowledge — if the user's question references their own documents, exhaust every discovery path first."""
AGENT_INSTRUCTIONS = "\n\n".join([
_INSTRUCTIONS_HEADER,
_READING_WORKFLOW,
_TOOL_USAGE_RULES,
_DISCOVERY,
_DECISION,
_AFTER_DISCOVERY,
_PERSISTENCE,
])
def _base_instructions(client) -> str:
"""Cloud: the live instructions the MCP server serves for this key's
tool set. Local: the built-in subset instructions."""
if not getattr(client, "api_key", None):
return AGENT_INSTRUCTIONS
instructions = _cloud_bridge(client).instructions()
if not isinstance(instructions, str) or not instructions.strip():
raise PageIndexAPIError(
"The MCP server returned no agent instructions — refusing to "
"substitute the SDK's local-subset guidance, which does not "
"cover the cloud tool set."
)
return instructions
def doc_targeting_block(client, doc_id, scoped: bool = False) -> Optional[str]:
"""The doc_id targeting text: names, metadata, and the directive to work
within those documents. Shared by agent_instructions and the local chat
surfaces (a leading conversation item on the OpenAI surfaces, a system
block on messages()). Raises when a doc_id's name is shadowed by a newer
same-name document — the name-addressed tools could not reach it. With
``scoped`` (surfaces whose tools resolve names inside the doc_id
allowlist) only a same-name duplicate within the targeted set
shadows."""
if doc_id is None:
return None
doc_ids = [doc_id] if isinstance(doc_id, str) else list(doc_id)
if not doc_ids:
return None
details = [client.get_document(one_id) for one_id in doc_ids]
listing = _all_documents(client)
documents = ([{**detail, "id": one_id}
for one_id, detail in zip(doc_ids, details)]
if scoped else listing)
for one_id, detail in zip(doc_ids, details):
entry, _ = _resolve_document(client, str(detail.get("name")),
documents=documents)
if entry is not None and entry.get("id") != one_id:
raise PageIndexAPIError(
f'Document "{detail.get("name")}" (doc_id: {one_id}) is '
"shadowed by a newer document with the same name (doc_id: "
f'{entry.get("id")}). The tools address documents by name '
"and would read the newer one. Rename or remove the "
"duplicate, or pass the newer doc_id."
)
# get_document keeps the cloud detail wire shape, which local mode
# serves without the user's metadata tags; the listing carries them
# in both modes.
by_id = {doc.get("id"): doc for doc in listing}
for one_id, detail in zip(doc_ids, details):
if detail.get("metadata") is None:
tags = _flat_metadata(by_id.get(one_id, {}).get("metadata"))
if tags is not None:
detail["metadata"] = tags
context = json.dumps(details, ensure_ascii=False)
if len(details) == 1:
return (
f"The user has specified document: {details[0].get('name')}\n"
f"Document metadata: {context}\n"
"Use this document's name to retrieve its content with "
"get_document_structure() and get_page_content()."
)
names = ", ".join(str(item.get("name")) for item in details)
return (
f"The user has specified documents: {names}\n"
f"Documents metadata: {context}\n"
"Use these documents' names to retrieve their content with "
"get_document_structure() and get_page_content()."
)
def build_agent_instructions(client, doc_id=None, scoped: bool = False) -> str:
"""Orchestration guidance for document QA agents; with doc_id, appends
the target documents and directs the agent to work within them."""
base = _base_instructions(client)
block = doc_targeting_block(client, doc_id, scoped=scoped)
return base if block is None else base + "\n\n" + block