Compare commits

..

1 Commits

Author SHA1 Message Date
dependabot[bot] fefaa7de80 build(deps): bump actions/setup-node from 5 to 6
Bumps [actions/setup-node](https://github.com/actions/setup-node) from 5 to 6.
- [Release notes](https://github.com/actions/setup-node/releases)
- [Commits](https://github.com/actions/setup-node/compare/v5...v6)

---
updated-dependencies:
- dependency-name: actions/setup-node
  dependency-version: '6'
  dependency-type: direct:production
  update-type: version-update:semver-major
...

Signed-off-by: dependabot[bot] <support@github.com>
2025-10-20 03:24:20 +00:00
24 changed files with 427 additions and 1474 deletions
+1 -1
View File
@@ -21,7 +21,7 @@ jobs:
- uses: pnpm/action-setup@v4
- name: Setup Node.js
uses: actions/setup-node@v5
uses: actions/setup-node@v6
with:
node-version-file: "ts/llama_cloud_services/.nvmrc"
+1 -1
View File
@@ -31,7 +31,7 @@ jobs:
- uses: pnpm/action-setup@v4
- name: Setup Node.js
uses: actions/setup-node@v5
uses: actions/setup-node@v6
with:
node-version-file: "ts/llama_cloud_services/.nvmrc"
- name: Install dependencies
+1 -1
View File
@@ -24,7 +24,7 @@ jobs:
- uses: actions/checkout@v5
- uses: pnpm/action-setup@v4
- name: Setup Node.js
uses: actions/setup-node@v5
uses: actions/setup-node@v6
with:
node-version-file: "ts/llama_cloud_services/.nvmrc"
- name: Install dependencies
@@ -20,7 +20,7 @@ jobs:
- uses: pnpm/action-setup@v4
- name: Setup Node.js
uses: actions/setup-node@v5
uses: actions/setup-node@v6
with:
node-version: "22"
cache: "pnpm"
@@ -4,19 +4,31 @@
"cell_type": "markdown",
"metadata": {},
"source": [
"# Document Classification + Extraction Workflow with LlamaCloud + LlamaIndex Workflows\n",
"# Complete Parse → Classify → Extract Workflow with LlamaCloud Services\n",
"\n",
"<a href=\"https://colab.research.google.com/github/run-llama/llama_cloud_services/blob/main/examples/misc/parse_classify_extract_workflow.ipynb\" target=\"_parent\"><img src=\"https://colab.research.google.com/assets/colab-badge.svg\" alt=\"Open In Colab\"/></a>\n",
"\n",
"This notebook shows a multi-step agentic document workflow that uses the **parsing**, **classification** and **extraction** modules in LlamaCloud, orchestrated through **LlamaIndex Workflows**. The workflow can take in a complex input document, parse it into clean markdown, classify it according to its subtype, and extract data according to a specified schema for that subtype. This allows you to automate document extraction of various types within the same workflow instead of having to manually separate the data beforehand. \n",
"\n",
"This notebook uses the following modules:\n",
"1. **Parse (LlamaParse)** - Extract and convert documents to markdown\n",
"This notebook demonstrates the complete workflow for processing documents using LlamaCloud services:\n",
"1. **Parse** - Extract and convert documents to markdown\n",
"2. **Classify** - Categorize documents based on their content\n",
"3. **Extract (LlamaExtract)** - Extract structured data using the markdown as input via SourceText\n",
"4. **LlamaIndex Workflows** - Event-driven orchestration of the parse, classify and extract steps\n",
"3. **Extract** - Extract structured data using the markdown as input via SourceText\n",
"\n",
"The workflow is implemented as a proper LlamaIndex Workflow with separate steps for parsing, classification, and extraction, connected by typed events. This provides modularity, observability, and type safety."
"## Overview of the Workflow\n",
"\n",
"### 1. Parse Phase\n",
"- Use `LlamaParse` to convert documents (PDFs, Word docs, etc.) into structured formats\n",
"- Extract markdown content that preserves document structure\n",
"- Get both raw text and markdown representations\n",
"\n",
"### 2. Classify Phase\n",
"- Use `ClassifyClient` to categorize documents based on content\n",
"- Apply classification rules to route documents appropriately\n",
"- Handle different document types with specific processing logic\n",
"\n",
"### 3. Extract Phase\n",
"- Use `LlamaExtract` with `SourceText` to extract structured data\n",
"- Pass the markdown content as input for more accurate extraction\n",
"- Define custom schemas for structured data extraction\n",
"\n",
"Let's walk through each step with practical examples."
]
},
{
@@ -33,8 +45,8 @@
"outputs": [],
"source": [
"# Install required packages\n",
"%pip install llama-cloud-services\n",
"%pip install python-dotenv"
"!pip install llama-cloud-services\n",
"!pip install python-dotenv"
]
},
{
@@ -61,7 +73,7 @@
"nest_asyncio.apply()\n",
"\n",
"# Set up API key\n",
"# os.environ[\"LLAMA_CLOUD_API_KEY\"] = \"\" # edit it\n",
"os.environ[\"LLAMA_CLOUD_API_KEY\"] = \"\" # edit it\n",
"\n",
"# Setup Base URL\n",
"# os.envrion[\"LLAMA_CLOUD_BASE_URL\"] = \"https://api.cloud.eu.llamaindex.ai/\" # update if necessay\n",
@@ -87,8 +99,7 @@
"name": "stdout",
"output_type": "stream",
"text": [
"Downloading financial_report.pdf...\n",
"✅ Downloaded financial_report.pdf\n",
"📁 financial_report.pdf already exists\n",
"📁 technical_spec.pdf already exists\n",
"\n",
"📂 Sample documents ready!\n"
@@ -104,7 +115,7 @@
"\n",
"# Download sample documents\n",
"docs_to_download = {\n",
" \"financial_report.pdf\": \"https://raw.githubusercontent.com/run-llama/llama_index/main/docs/examples/data/10k/uber_2021.pdf\",\n",
" \"financial_report.pdf\": \"https://raw.githubusercontent.com/run-llama/llama_index/main/docs/docs/examples/data/10k/uber_2021.pdf\",\n",
" \"technical_spec.pdf\": \"https://www.ti.com/lit/ds/symlink/lm317.pdf\",\n",
"}\n",
"\n",
@@ -144,10 +155,10 @@
"output_type": "stream",
"text": [
"🔄 Parsing documents...\n",
"Started parsing the file under job_id 530c187a-bd2d-4eea-b38d-9e5738eab465\n",
".✅ Parsed financial report (Job ID: 530c187a-bd2d-4eea-b38d-9e5738eab465)\n",
"Started parsing the file under job_id a6e27710-776b-4445-8b94-8d75959ff5db\n",
"✅ Parsed technical spec (Job ID: a6e27710-776b-4445-8b94-8d75959ff5db)\n",
"Started parsing the file under job_id 8a8c76f9-354d-4275-91d8-312ff1adc762\n",
"...✅ Parsed financial report (Job ID: 8a8c76f9-354d-4275-91d8-312ff1adc762)\n",
"Started parsing the file under job_id 7e603448-ed80-4d18-948b-6801ed51c41b\n",
"✅ Parsed technical spec (Job ID: 7e603448-ed80-4d18-948b-6801ed51c41b)\n",
"\n",
"📄 Parsing complete!\n"
]
@@ -235,23 +246,23 @@
"\n",
"## 1 Features\n",
"\n",
"- Output voltage range:\n",
" Output voltage range:\n",
" Adjustable: 1.25V to 37V\n",
"- Output current: 1.5A\n",
"- Line regulation: 0.01%/V (typ)\n",
"- Load regulation: 0.1% (typ)\n",
"- Internal short-circuit current limiting\n",
"- Thermal overload protection\n",
"- Output safe-area compensation (new chip)\n",
"- PSRR: 80dB at 120Hz for CADJ = 10μF (new chip)\n",
"- Packages:\n",
" Output current: 1.5A\n",
" Line regulation: 0.01%/V (typ)\n",
" Load regulation: 0.1% (typ)\n",
" Internal short-circuit current limiting\n",
" Thermal overload protection\n",
" Output safe-area compensation (new chip)\n",
" PSRR: 80dB at 120Hz for CADJ = 10μF (new chip)\n",
" Packages:\n",
" 4-pin, SOT-223 (DCY)\n",
" 3-pin, TO-263 (KTT)\n",
" 3-pin, TO-220 (KCS, KCT),\n",
"...\n",
"\n",
"📏 Financial report markdown length: 1338499 characters\n",
"📏 Technical spec markdown length: 92483 characters\n"
"📏 Financial report markdown length: 1348671 characters\n",
"📏 Technical spec markdown length: 90971 characters\n"
]
}
],
@@ -328,72 +339,6 @@
"print(f\"📝 Created {len(classification_rules)} classification rules\")"
]
},
{
"cell_type": "markdown",
"metadata": {},
"source": [
"### Try Classification Independently\n",
"\n",
"Let's test the classification on one of our parsed documents to see how it works:\n"
]
},
{
"cell_type": "code",
"execution_count": null,
"metadata": {},
"outputs": [
{
"name": "stdout",
"output_type": "stream",
"text": [
"🔍 Classifying financial document...\n",
" Document length: 1,338,499 characters\n",
"\n",
"✅ Classification Result:\n",
" Type: financial_document\n",
" Confidence: 100.00%\n",
" Reasoning: This document is a Form 10-K, which is an annual report required by the U.S. Securities and Exchange Commission (SEC) for publicly traded companies. It contains financial data, information about the c...\n",
"\n",
"======================================================================\n"
]
}
],
"source": [
"# Let's classify the financial document\n",
"print(\"🔍 Classifying financial document...\")\n",
"print(f\" Document length: {len(financial_markdown):,} characters\\n\")\n",
"\n",
"# Write to temp file for classification\n",
"import tempfile\n",
"from pathlib import Path\n",
"\n",
"with tempfile.NamedTemporaryFile(\n",
" mode=\"w\", suffix=\".md\", delete=False, encoding=\"utf-8\"\n",
") as tmp:\n",
" tmp.write(financial_markdown)\n",
" temp_financial_path = Path(tmp.name)\n",
"\n",
"# Classify the document\n",
"financial_classification = await classify_client.aclassify_file_path(\n",
" rules=classification_rules, file_input_path=str(temp_financial_path)\n",
")\n",
"\n",
"doc_type = financial_classification.items[0].result.type\n",
"confidence = financial_classification.items[0].result.confidence\n",
"reasoning = financial_classification.items[0].result.reasoning\n",
"\n",
"print(f\"✅ Classification Result:\")\n",
"print(f\" Type: {doc_type}\")\n",
"print(f\" Confidence: {confidence:.2%}\")\n",
"print(\n",
" f\" Reasoning: {reasoning[:200]}...\"\n",
" if reasoning and len(reasoning) > 200\n",
" else f\" Reasoning: {reasoning}\"\n",
")\n",
"\n",
"print(\"\\n\" + \"=\" * 70)"
]
},
{
"cell_type": "markdown",
"metadata": {},
@@ -499,31 +444,9 @@
"cell_type": "markdown",
"metadata": {},
"source": [
"## Building the Complete Workflow\n",
"## Complete Workflow Summary\n",
"\n",
"Now that we've seen how parsing works, let's build a complete 3-step workflow (Parse → Classify → Extract) using LlamaIndex Workflows. We'll define the workflow structure here, and you can see it in action below where we also demonstrate the classification and extraction modules independently.\n",
"\n",
"### Install Workflows Package\n",
"\n",
"First, let's install the LlamaIndex workflows package:\n"
]
},
{
"cell_type": "code",
"execution_count": null,
"metadata": {},
"outputs": [],
"source": [
"%pip install llama-index-workflows llama-index-utils-workflow"
]
},
{
"cell_type": "markdown",
"metadata": {},
"source": [
"## Define the Workflow\n",
"\n",
"Let's restructure the document processing into a proper LlamaIndex Workflow with separate classification and extraction steps:\n"
"Let's create a function that demonstrates the complete workflow:"
]
},
{
@@ -535,7 +458,7 @@
"name": "stdout",
"output_type": "stream",
"text": [
"🔧 Workflow defined!\n"
"🔧 Workflow function defined!\n"
]
}
],
@@ -543,286 +466,81 @@
"import tempfile\n",
"from pathlib import Path\n",
"from llama_cloud import ExtractConfig\n",
"from workflows import Workflow, step, Context\n",
"from workflows.events import Event, StartEvent, StopEvent\n",
"\n",
"\n",
"# Define workflow events\n",
"class ParseEvent(Event):\n",
" \"\"\"Event emitted after parsing\"\"\"\n",
"\n",
" file_path: str\n",
" markdown_content: str\n",
" job_id: str\n",
"\n",
"\n",
"class ClassifyEvent(Event):\n",
" \"\"\"Event emitted after classification\"\"\"\n",
"\n",
" markdown_content: str\n",
" temp_path: str\n",
" doc_type: str\n",
" confidence: float\n",
"\n",
"\n",
"class ExtractEvent(Event):\n",
" \"\"\"Event emitted after extraction\"\"\"\n",
"\n",
" doc_type: str\n",
" confidence: float\n",
" extracted_data: dict\n",
" markdown_length: int\n",
" temp_path: str\n",
" markdown_sample: str\n",
"\n",
"\n",
"class DocumentWorkflow(Workflow):\n",
"async def complete_document_workflow(markdown_content: str):\n",
" \"\"\"\n",
" Complete document processing workflow: Parse → Classify → Extract\n",
" Complete workflow: Parse → Classify → Extract\n",
" \"\"\"\n",
" print(f\"🚀 Starting complete workflow\")\n",
" print(\"=\" * 60)\n",
"\n",
" def __init__(\n",
" self,\n",
" parser,\n",
" classify_client,\n",
" classification_rules,\n",
" llama_extract,\n",
" financial_schema,\n",
" technical_schema,\n",
" **kwargs,\n",
" ):\n",
" super().__init__(**kwargs)\n",
" self.parser = parser\n",
" self.classify_client = classify_client\n",
" self.classification_rules = classification_rules\n",
" self.llama_extract = llama_extract\n",
" self.financial_schema = financial_schema\n",
" self.technical_schema = technical_schema\n",
" # Step 1: Classify\n",
" print(\"🏷️ Step 2: Classifying document...\")\n",
"\n",
" @step\n",
" async def parse_document(self, ctx: Context, ev: StartEvent) -> ParseEvent:\n",
" \"\"\"\n",
" Step 1: Parse the document to extract markdown\n",
" \"\"\"\n",
" file_path = ev.file_path\n",
" print(f\"📄 Step 1: Parsing document: {file_path}...\")\n",
" with tempfile.NamedTemporaryFile(\n",
" mode=\"w\", suffix=\".md\", delete=False, encoding=\"utf-8\"\n",
" ) as tmp:\n",
" tmp.write(markdown_content)\n",
" temp_path = Path(tmp.name)\n",
"\n",
" # Parse the document\n",
" parse_result = await self.parser.aparse(file_path)\n",
" markdown_content = await parse_result.aget_markdown()\n",
" job_id = parse_result.job_id\n",
" print(temp_path)\n",
"\n",
" print(f\" ✅ Parsed successfully (Job ID: {job_id})\")\n",
" print(f\" 📝 Extracted {len(markdown_content):,} characters\")\n",
" classification = await classify_client.aclassify_file_path(\n",
" rules=classification_rules, file_input_path=str(temp_path)\n",
" )\n",
" doc_type = classification.items[0].result.type\n",
" confidence = classification.items[0].result.confidence\n",
" print(f\" ✅ Classified as: {doc_type} (confidence: {confidence:.2f})\")\n",
"\n",
" # Write event to stream for monitoring\n",
" parse_event = ParseEvent(\n",
" file_path=file_path,\n",
" markdown_content=markdown_content,\n",
" job_id=job_id,\n",
" )\n",
" ctx.write_event_to_stream(parse_event)\n",
" # Step 2: Extract based on classification\n",
" print(\"🔍 Step 3: Extracting structured data using SourceText...\")\n",
" source_text = SourceText(\n",
" text_content=markdown_content,\n",
" filename=f\"{os.path.basename(temp_path)}_markdown.md\",\n",
" )\n",
"\n",
" return parse_event\n",
" # Choose schema based on classification\n",
" if \"financial\" in doc_type.lower():\n",
" schema = FinancialMetrics\n",
" print(\" 📊 Using FinancialMetrics schema\")\n",
" elif \"technical\" in doc_type.lower():\n",
" schema = TechnicalSpec\n",
" print(\" 🔧 Using TechnicalSpec schema\")\n",
" else:\n",
" schema = FinancialMetrics # Default fallback\n",
" print(\" 📊 Using default FinancialMetrics schema\")\n",
"\n",
" @step\n",
" async def classify_document(self, ctx: Context, ev: ParseEvent) -> ClassifyEvent:\n",
" \"\"\"\n",
" Step 2: Classify the document based on its content\n",
" \"\"\"\n",
" markdown_content = ev.markdown_content\n",
" print(\"🏷️ Step 2: Classifying document...\")\n",
" extract_config = ExtractConfig(\n",
" extraction_mode=\"BALANCED\",\n",
" )\n",
"\n",
" # Write markdown to temp file for classification\n",
" with tempfile.NamedTemporaryFile(\n",
" mode=\"w\", suffix=\".md\", delete=False, encoding=\"utf-8\"\n",
" ) as tmp:\n",
" tmp.write(markdown_content)\n",
" temp_path = Path(tmp.name)\n",
" extraction_result = llama_extract.extract(\n",
" data_schema=schema, config=extract_config, files=source_text\n",
" )\n",
"\n",
" # Classify the document\n",
" classification = await self.classify_client.aclassify_file_path(\n",
" rules=self.classification_rules, file_input_path=str(temp_path)\n",
" )\n",
" doc_type = classification.items[0].result.type\n",
" confidence = classification.items[0].result.confidence\n",
" print(\" ✅ Extraction complete!\")\n",
"\n",
" print(f\" ✅ Classified as: {doc_type} (confidence: {confidence:.2f})\")\n",
"\n",
" # Write event to stream for monitoring\n",
" classify_event = ClassifyEvent(\n",
" markdown_content=markdown_content,\n",
" temp_path=str(temp_path),\n",
" doc_type=doc_type,\n",
" confidence=confidence,\n",
" )\n",
" ctx.write_event_to_stream(classify_event)\n",
"\n",
" return classify_event\n",
"\n",
" @step\n",
" async def extract_data(self, ctx: Context, ev: ClassifyEvent) -> ExtractEvent:\n",
" \"\"\"\n",
" Step 3: Extract structured data based on classification\n",
" \"\"\"\n",
" print(\"🔍 Step 3: Extracting structured data using SourceText...\")\n",
"\n",
" # Choose schema based on classification\n",
" if \"financial\" in ev.doc_type.lower():\n",
" schema = self.financial_schema\n",
" print(\" 📊 Using FinancialMetrics schema\")\n",
" elif \"technical\" in ev.doc_type.lower():\n",
" schema = self.technical_schema\n",
" print(\" 🔧 Using TechnicalSpec schema\")\n",
" else:\n",
" schema = self.financial_schema # Default fallback\n",
" print(\" 📊 Using default FinancialMetrics schema\")\n",
"\n",
" # Create SourceText from markdown content\n",
" source_text = SourceText(\n",
" text_content=ev.markdown_content,\n",
" filename=f\"{os.path.basename(ev.temp_path)}_markdown.md\",\n",
" )\n",
"\n",
" # Configure extraction\n",
" extract_config = ExtractConfig(\n",
" extraction_mode=\"BALANCED\",\n",
" )\n",
"\n",
" # Perform extraction\n",
" extraction_result = self.llama_extract.extract(\n",
" data_schema=schema, config=extract_config, files=source_text\n",
" )\n",
"\n",
" print(\" ✅ Extraction complete!\")\n",
"\n",
" # Create markdown sample\n",
" markdown_sample = (\n",
" ev.markdown_content[:200] + \"...\"\n",
" if len(ev.markdown_content) > 200\n",
" else ev.markdown_content\n",
" )\n",
"\n",
" extract_event = ExtractEvent(\n",
" doc_type=ev.doc_type,\n",
" confidence=ev.confidence,\n",
" extracted_data=extraction_result.data,\n",
" markdown_length=len(ev.markdown_content),\n",
" temp_path=ev.temp_path,\n",
" markdown_sample=markdown_sample,\n",
" )\n",
" ctx.write_event_to_stream(extract_event)\n",
"\n",
" return extract_event\n",
"\n",
" @step\n",
" async def finalize_results(self, ctx: Context, ev: ExtractEvent) -> StopEvent:\n",
" \"\"\"\n",
" Step 4: Finalize and return results\n",
" \"\"\"\n",
" result = {\n",
" \"file_path\": ev.temp_path,\n",
" \"markdown_length\": ev.markdown_length,\n",
" \"classification\": ev.doc_type,\n",
" \"confidence\": ev.confidence,\n",
" \"extracted_data\": ev.extracted_data,\n",
" \"markdown_sample\": ev.markdown_sample,\n",
" }\n",
"\n",
" return StopEvent(result=result)\n",
" return {\n",
" \"file_path\": temp_path,\n",
" \"markdown_length\": len(markdown_content),\n",
" \"classification\": doc_type,\n",
" \"confidence\": confidence,\n",
" \"extracted_data\": extraction_result.data,\n",
" \"markdown_sample\": markdown_content[:200] + \"...\"\n",
" if len(markdown_content) > 200\n",
" else markdown_content,\n",
" }\n",
"\n",
"\n",
"print(\"🔧 Workflow defined!\")"
"print(\"🔧 Workflow function defined!\")"
]
},
{
"cell_type": "markdown",
"metadata": {},
"source": [
"### Workflow Structure\n",
"\n",
"The workflow consists of four steps connected by typed events:\n",
"\n",
"```\n",
"┌─────────────┐\n",
"│ StartEvent │ (file_path)\n",
"└──────┬──────┘\n",
" │\n",
" ▼\n",
"┌──────────────────┐\n",
"│ parse_document │ Step 1: Parse PDF to markdown\n",
"└──────┬───────────┘\n",
" │\n",
" ▼\n",
"┌─────────────┐\n",
"│ ParseEvent │ (markdown_content, job_id)\n",
"└──────┬──────┘\n",
" │\n",
" ▼\n",
"┌─────────────────────┐\n",
"│ classify_document │ Step 2: Classification\n",
"└──────┬──────────────┘\n",
" │\n",
" ▼\n",
"┌──────────────┐\n",
"│ ClassifyEvent│ (doc_type, confidence, markdown_content)\n",
"└──────┬───────┘\n",
" │\n",
" ▼\n",
"┌──────────────┐\n",
"│ extract_data │ Step 3: Extraction with schema selection\n",
"└──────┬───────┘\n",
" │\n",
" ▼\n",
"┌──────────────┐\n",
"│ ExtractEvent │ (extracted_data, doc_type, confidence)\n",
"└──────┬───────┘\n",
" │\n",
" ▼\n",
"┌──────────────────┐\n",
"│ finalize_results │ Step 4: Format and return results\n",
"└──────┬───────────┘\n",
" │\n",
" ▼\n",
"┌─────────────┐\n",
"│ StopEvent │ (final result dictionary)\n",
"└─────────────┘\n",
"```\n",
"\n",
"**Key Features:**\n",
"- **Step 1 (parse_document)**: Takes a file path and parses the document into clean markdown\n",
"- **Step 2 (classify_document)**: Takes markdown content and classifies it into document types\n",
"- **Step 3 (extract_data)**: Selects appropriate schema based on classification and extracts structured data\n",
"- **Step 4 (finalize_results)**: Packages all results into final output format\n",
"- Events are written to the stream for real-time monitoring\n"
]
},
{
"cell_type": "markdown",
"metadata": {},
"source": [
"## Visualize the Workflow\n",
"\n",
"Let's visualize the workflow structure to see the flow of events:\n"
]
},
{
"cell_type": "code",
"execution_count": null,
"metadata": {},
"outputs": [],
"source": [
"# Initialize the workflow\n",
"workflow = DocumentWorkflow(\n",
" parser=parser,\n",
" classify_client=classify_client,\n",
" classification_rules=classification_rules,\n",
" llama_extract=llama_extract,\n",
" financial_schema=FinancialMetrics,\n",
" technical_schema=TechnicalSpec,\n",
" timeout=300,\n",
" verbose=True,\n",
")"
"## Run Complete Workflow on Both Documents"
]
},
{
@@ -834,173 +552,53 @@
"name": "stdout",
"output_type": "stream",
"text": [
"document_workflow.html\n"
]
}
],
"source": [
"# Draw the workflow visualization\n",
"from llama_index.utils.workflow import draw_all_possible_flows\n",
"\n",
"draw_all_possible_flows(\n",
" workflow,\n",
" filename=\"document_workflow.html\",\n",
")"
]
},
{
"cell_type": "markdown",
"metadata": {},
"source": [
"The workflow has been visualized and saved to `document_workflow.html`. You can open this file in a browser to see the interactive workflow diagram.\n"
]
},
{
"cell_type": "markdown",
"metadata": {},
"source": [
"The workflow visualization shows:\n",
"1. **StartEvent** → **parse_document** step\n",
"2. **ParseEvent** → **classify_document** step\n",
"3. **ClassifyEvent** → **extract_data** step \n",
"4. **ExtractEvent** → **finalize_results** step\n",
"5. **StopEvent** (final output)\n",
"\n",
"Each step is connected by typed events, allowing for clean separation of concerns and easy monitoring of the workflow execution.\n"
]
},
{
"cell_type": "markdown",
"metadata": {},
"source": [
"## Run the Workflow on Both Documents\n",
"\n",
"Now let's run the workflow on both documents and monitor the events:\n"
]
},
{
"cell_type": "code",
"execution_count": null,
"metadata": {},
"outputs": [
{
"name": "stdout",
"output_type": "stream",
"text": [
"\n",
"======================================================================\n",
"🚀 Processing Document 1: sample_docs/financial_report.pdf\n",
"======================================================================\n",
"\n",
"Running step parse_document\n",
"📄 Step 1: Parsing document: sample_docs/financial_report.pdf...\n",
"Started parsing the file under job_id bb53c6bf-79cc-4f63-9c97-16983d59f29d\n",
". ✅ Parsed successfully (Job ID: bb53c6bf-79cc-4f63-9c97-16983d59f29d)\n",
" 📝 Extracted 1,338,499 characters\n",
"Step parse_document produced event ParseEvent\n",
"📄 Parse Event: Extracted 1,338,499 characters\n",
"Running step classify_document\n",
"🚀 Starting complete workflow\n",
"============================================================\n",
"🏷️ Step 2: Classifying document...\n",
"/var/folders/g6/4b5lpp5974gcpr890ybhbw4r0000gn/T/tmpos3b62tm.md\n",
" ✅ Classified as: financial_document (confidence: 1.00)\n",
"Step classify_document produced event ClassifyEvent\n",
"📊 Classification Event: financial_document (1.00)\n",
"Running step extract_data\n",
"🔍 Step 3: Extracting structured data using SourceText...\n",
" 📊 Using FinancialMetrics schema\n",
".. ✅ Extraction complete!\n",
"Step extract_data produced event ExtractEvent\n",
"Running step finalize_results\n",
"Step finalize_results produced event StopEvent\n",
"✅ Extraction Event: 7 fields extracted\n",
"\n",
"✅ Document 1 processed successfully!\n",
"============================================================\n",
"\n",
"======================================================================\n",
"🚀 Processing Document 2: sample_docs/technical_spec.pdf\n",
"======================================================================\n",
"\n",
"Running step parse_document\n",
"📄 Step 1: Parsing document: sample_docs/technical_spec.pdf...\n",
"Started parsing the file under job_id 944905c1-3c49-431a-ad86-4436d16f3d1c\n",
" ✅ Parsed successfully (Job ID: 944905c1-3c49-431a-ad86-4436d16f3d1c)\n",
" 📝 Extracted 92,483 characters\n",
"Step parse_document produced event ParseEvent\n",
"📄 Parse Event: Extracted 92,483 characters\n",
"Running step classify_document\n",
"🚀 Starting complete workflow\n",
"============================================================\n",
"🏷️ Step 2: Classifying document...\n",
"/var/folders/g6/4b5lpp5974gcpr890ybhbw4r0000gn/T/tmpppz9ub_m.md\n",
" ✅ Classified as: technical_specification (confidence: 1.00)\n",
"Step classify_document produced event ClassifyEvent\n",
"📊 Classification Event: technical_specification (1.00)\n",
"Running step extract_data\n",
"🔍 Step 3: Extracting structured data using SourceText...\n",
" 🔧 Using TechnicalSpec schema\n",
" ✅ Extraction complete!\n",
"Step extract_data produced event ExtractEvent\n",
"Running step finalize_results\n",
"Step finalize_results produced event StopEvent\n",
"✅ Extraction Event: 8 fields extracted\n",
"\n",
"✅ Document 2 processed successfully!\n",
"\n",
"============================================================\n",
"\n",
"📋 Processed 2 documents successfully!\n"
]
}
],
"source": [
"# Process both documents through the workflow\n",
"# Process both documents through the complete workflow\n",
"results = []\n",
"\n",
"# Define the document files to process\n",
"document_files = [\n",
" \"sample_docs/financial_report.pdf\",\n",
" \"sample_docs/technical_spec.pdf\",\n",
"]\n",
"\n",
"for i, file_path in enumerate(document_files, 1):\n",
" print(f\"\\n{'='*70}\")\n",
" print(f\"🚀 Processing Document {i}: {file_path}\")\n",
" print(f\"{'='*70}\\n\")\n",
"\n",
"for doc_text in document_texts:\n",
" try:\n",
" # Run the workflow\n",
" handler = workflow.run(file_path=file_path)\n",
"\n",
" # Monitor events as they are emitted\n",
" async for event in handler.stream_events():\n",
" if isinstance(event, ParseEvent):\n",
" print(\n",
" f\"📄 Parse Event: Extracted {len(event.markdown_content):,} characters\"\n",
" )\n",
" elif isinstance(event, ClassifyEvent):\n",
" print(\n",
" f\"📊 Classification Event: {event.doc_type} ({event.confidence:.2f})\"\n",
" )\n",
" elif isinstance(event, ExtractEvent):\n",
" print(\n",
" f\"✅ Extraction Event: {len(event.extracted_data)} fields extracted\"\n",
" )\n",
"\n",
" # Get final result\n",
" result = await handler\n",
" result = await complete_document_workflow(doc_text)\n",
" results.append(result)\n",
"\n",
" print(f\"\\n✅ Document {i} processed successfully!\")\n",
"\n",
" print(\"\\n\" + \"=\" * 60 + \"\\n\")\n",
" except Exception as e:\n",
" print(f\"❌ Error processing document {i}: {str(e)}\")\n",
" import traceback\n",
" print(f\"❌ Error processing {doc_path}: {str(e)}\")\n",
" print(\"\\n\" + \"=\" * 60 + \"\\n\")\n",
"\n",
" traceback.print_exc()\n",
"\n",
"print(f\"\\n\\n📋 Processed {len(results)} documents successfully!\")"
"print(f\"📋 Processed {len(results)} documents successfully!\")"
]
},
{
"cell_type": "markdown",
"metadata": {},
"source": [
"## Final Results Summary\n"
"## Final Results Summary"
]
},
{
@@ -1015,9 +613,9 @@
"📈 COMPLETE WORKFLOW RESULTS SUMMARY\n",
"======================================================================\n",
"\n",
"📄 Document 1: tmpuyxzpd3x.md\n",
"📄 Document 1: tmpos3b62tm.md\n",
" 📊 Classification: financial_document (confidence: 1.00)\n",
" 📝 Markdown length: 1,338,499 characters\n",
" 📝 Markdown length: 1,348,671 characters\n",
" 📋 Markdown sample: \n",
"\n",
"# UNITED STATES\n",
@@ -1031,14 +629,14 @@
" • company_name: Uber Technologies, Inc.\n",
" • document_type: Annual Report on Form 10-K\n",
" • fiscal_year: 2021\n",
" • revenue_2021: $17,455 and $21,764\n",
" • net_income_2021: $(496) to (700)\n",
" • key_business_segments: ['Borrower and the Restricted Subsidiaries', 'Holdings', 'Guarantors', 'Material Domestic Subsidiaries', 'Material Foreign Subsidiaries']\n",
" • risk_factors: ['Indemnification obligations of the borrower for losses, claims, damages, liabilities, and out-of-pocket expenses incurred by agents, lenders, arrangers, and related parties in connection with the agreement or loans, except in certain cases such as gross negligence, bad faith, willful misconduct, or material breach by the indemnitee.', \"Borrower not required to indemnify any indemnitee for settlements entered into without the borrower's consent.\", 'Limitation of liability for special, indirect, consequential, or punitive damages, and for damages from unauthorized use of information, except for direct damages resulting from gross negligence, bad faith, or willful misconduct.', 'Obligation of the borrower to indemnify the administrative agent for liabilities arising from performance of duties, except in cases of gross negligence, bad faith, or willful misconduct.', 'Limitations and conditions on assignments and participations of lender rights, including restrictions on assignments to disqualified institutions, loan parties, affiliates of loan parties, defaulting lenders, and natural persons.', 'Setoff rights for lenders and issuing banks after an event of default, allowing them to apply borrower deposits toward obligations under the agreement.', 'Potential for increased obligations under the agreement as a result of changes in law affecting payment terms.', 'Requirement for the borrower and guarantors to provide information to comply with anti-money laundering rules and the USA PATRIOT Act.']\n",
" • revenue_2021: $21,764\n",
" • net_income_2021: $(496)\n",
" • key_business_segments: ['Mobility', 'Delivery', 'Freight', 'All Other (including former New Mobility, e-bikes, e-scooters, Advanced Technologies Group and other technology programs)']\n",
" • risk_factors: [\"The company faces numerous risk factors across its business operations and environment. The COVID-19 pandemic and related mitigation measures have adversely affected parts of the business, including reduced demand for Mobility offerings and creating ongoing uncertainties. The company's operational and financial performance is influenced by competitive pressure in the mobility, delivery, and logistics industries, characterized by well-established alternatives, low barriers to entry, and low switching costs. Driver classification risks exist if Drivers are deemed employees, workers, or quasi-employees rather than independent contractors, exposing the company to legal actions and financial liabilities globally. Competition challenges require the company to sometimes lower fares, offer incentives, and promotions, which impacts profitability. There are significant operating losses historically with substantial future operating expense increases anticipated, and the ability to achieve or maintain profitability is uncertain. Network value depends on maintaining critical mass among Drivers, consumers, merchants, shippers, and carriers, and failures to do so diminish platform attractiveness. Brand and reputation maintenance is critical, with exposure to negative publicity, media coverage, and risks from associated companies' brands or licensed brands in joint ventures.\\n\\nOperational risks include historical workplace culture and compliance challenges, management complexity due to rapid growth, technological infrastructure issues potentially causing disruptions or poor user experience, and security or data privacy breaches that could impact revenue and reputation. Platform users may engage in or be subjected to criminal, violent, or dangerous activity leading to safety incidents and legal actions. New offerings and technologies investments are inherently risky without guaranteed benefits. Economic conditions, inflation, and increased costs (fuel, food, labor, energy) may negatively impact results. Regulatory risks are extensive and global, involving payment and financial services compliance, licensing, anti-money laundering laws, data privacy (GDPR, CCPA, LGPD), and labor laws. Legal and regulatory investigations and inquiries, including antitrust, FCPA, labor classification, data protection, and intellectual property matters, pose risks of fines, penalties, operational changes, and increased costs.\\n\\nGeopolitical and jurisdictional risks include operating limitations or bans in some locations, currency exchange risk, and complex evolving regulations with the potential for fines and loss of licenses or permits. Insurance risks include potential inadequacy of reserves, liability exposure from accidents or impersonation, and insurer insolvency. Driver qualification requirements and background checks may increase costs or fail to expose all relevant information, with associated insurance cost risks and potential for courtroom or regulatory challenges to pricing models.\\n\\nFinancial risks comprise significant accumulated deficits, requirement for additional capital with uncertain availability, debt obligations, tax exposure including uncertain positions and observed changes in tax laws, and volatility in common stock price with no expected cash dividends. Accounting judgments and estimates involve critical assumptions affecting reported financial metrics related to goodwill, revenue recognition, incentive accruals, and stock-based compensation. Cybersecurity risks include exposures to malware, ransomware, phishing, and other cyberattacks. Climate change presents physical and transitional risks that may impact operations and costs, and failure to meet climate commitments may have operational and reputational consequences.\\n\\nOther risks include potential liability under anti-corruption and anti-terrorism laws, adverse effects from defaults under debt agreements, limitations in takeover actions due to corporate governance provisions, and the impact of non-GAAP financial measure limitations. Overall, these diverse and interconnected risk factors contribute to significant uncertainty regarding the company's future business prospects, operating results, and financial condition.\"]\n",
"\n",
"📄 Document 2: tmp7ower2xm.md\n",
"📄 Document 2: tmpppz9ub_m.md\n",
" 📊 Classification: technical_specification (confidence: 1.00)\n",
" 📝 Markdown length: 92,483 characters\n",
" 📝 Markdown length: 90,971 characters\n",
" 📋 Markdown sample: \n",
"\n",
"LM317\n",
@@ -1050,14 +648,20 @@
" 🎯 Extracted fields: 8 fields\n",
" • component_name: LM317\n",
" • manufacturer: Texas Instruments\n",
" • part_number: LM317, SLVS044Z\n",
" • description: The LM317 is an adjustable three-pin, positive-voltage regulator capable of supplying more than 1.5A (typically up to 1.5A) over an output voltage range of 1.25V to 37V. The device requires only two external resistors to set the output voltage. It features a typical line regulation of 0.01% and typical load regulation of 0.1%. The LM317 includes current limiting, thermal overload protection, and safe operating area protection. Overload protection remains functional even if the ADJUST pin is disconnected. The regulator is used in applications such as constant-current battery-charger circuits, slow turn-on 15V regulator circuits, AC voltage-regulator circuits, current-limited charger circuits, and high-current and adjustable regulator circuits. It is available in packages including SOT-223 (DCY), TO-220 (KCS), and TO-263 (KTT).\n",
" • part_number: LM317\n",
" • description: The LM317 is an adjustable three-pin, positive-voltage regulator capable of supplying up to 1.5A over an output voltage range of 1.25V to 37V. It features line and load regulation, internal current limiting, thermal overload protection, and safe operating area compensation.\n",
" • operating_voltage: {'min_voltage': 1.25, 'max_voltage': 37.0, 'unit': 'V'}\n",
" • maximum_current: 4.0\n",
" • key_features: ['Adjustable output voltage range: 1.25V to 37V', 'Output current up to 1.5A (up to 4A with external pass elements)', 'Line regulation: typically 0.01%/V', 'Load regulation: typically 0.1%', 'Internal short-circuit current limiting / Current limiting', 'Thermal overload protection / Thermal shutdown', 'Output safe-area compensation / Safe operating area protection', 'PSRR: 80dB at 120Hz for CADJ = 10μF (new chip)', 'NPN Darlington output drive', 'Programmable feedback', 'Multiple package options (SOT-223, TO-220, TO-263)', 'Can be used in constant-current, battery-charging, and regulator applications']\n",
" • applications: ['Multifunction printers, AC drive power stage modules, Electricity meters, Servo drive control modules, Merchant network and server PSU, Adjustable voltage regulator, 0V to 30V regulator circuit, Regulator circuit with improved ripple rejection, Precision current-limiter, Tracking preregulator, 1.25V to 20V regulator, Battery charger circuit, Constant-current battery charger circuits, Slow turn-on regulator, AC voltage-regulator, Current-limited charger circuits, High-current adjustable regulator circuits, General-purpose adjustable power supply']\n",
" • maximum_current: 1.5\n",
" • key_features: ['Adjustable output voltage: 1.25V to 37V', 'Output current up to 1.5A', 'Line regulation: 0.01%/V (typical)', 'Load regulation: 0.1% (typical)', 'Internal short-circuit current limiting', 'Thermal overload protection', 'Output safe-area compensation', 'High power supply rejection ratio (PSRR): 80dB at 120Hz (new chip)', 'Available in SOT-223, TO-263, and TO-220 packages']\n",
" • applications: ['Multifunction printers', 'AC drive power stage modules', 'Electricity meters', 'Servo drive control modules', 'Merchant network and server power supply units']\n",
"\n",
"✨ Workflow completed successfully!\n"
"✨ Workflow completed successfully!\n",
"\n",
"📚 Key Learnings:\n",
" • Parse: Converted documents to clean markdown format\n",
" • Classify: Automatically categorized document types\n",
" • Extract: Used SourceText with markdown for structured data extraction\n",
" • The markdown content provides much better context for extraction than raw PDFs\n"
]
}
],
@@ -1079,7 +683,14 @@
" for key, value in extracted.items():\n",
" print(f\" • {key}: {value}\")\n",
"\n",
"print(\"\\n✨ Workflow completed successfully!\")"
"print(\"\\n✨ Workflow completed successfully!\")\n",
"print(\"\\n📚 Key Learnings:\")\n",
"print(\" • Parse: Converted documents to clean markdown format\")\n",
"print(\" • Classify: Automatically categorized document types\")\n",
"print(\" • Extract: Used SourceText with markdown for structured data extraction\")\n",
"print(\n",
" \" • The markdown content provides much better context for extraction than raw PDFs\"\n",
")"
]
},
{
@@ -1088,9 +699,9 @@
"source": [
"## Conclusion\n",
"\n",
"The notebook shows you how to build an e2e document **Classify → Extract** workflow using LlamaCloud. This uses some of our core building blocks around **classification** interleaved with **document extraction**.\n",
"This notebook demonstrated the complete **Parse → Classify → Extract** workflow using LlamaCloud services:\n",
"\n",
"### Main Components:\n",
"### Key Components:\n",
"\n",
"1. **LlamaParse** (`llama_cloud_services.parse.base.LlamaParse`):\n",
" - Converts documents to clean, structured markdown\n",
@@ -1104,17 +715,38 @@
"\n",
"3. **LlamaExtract with SourceText** (`llama_cloud_services.extract.extract.LlamaExtract`, `SourceText`):\n",
" - Extracts structured data using custom Pydantic schemas\n",
" - You can either feed in the file directly (in which case parsing will happen under the hood), or the parsed text through the **SourceText** object (which is the case in this example) \n",
" - **SourceText** allows using markdown content as input instead of raw files\n",
" - Provides much better extraction accuracy when using processed markdown\n",
"\n",
"**Benefits of an e2e workflow**: The main benefit of doing Classify -> Extract, instead of only Extract, is the fact that you can handle documents of different types/different expected schemas within the same workflow, without having to separate out the data before and running separate extractions on each data subset. "
"### Workflow Benefits:\n",
"\n",
"- **Better Accuracy**: Using markdown from parsing provides cleaner, more structured input for extraction\n",
"- **Automatic Routing**: Classification allows different processing logic for different document types\n",
"- **Structured Output**: Custom schemas ensure consistent, structured data extraction\n",
"- **Flexible Input**: SourceText supports text content, file paths, and bytes\n",
"\n",
"### Key Insights:\n",
"\n",
"1. **SourceText is the bridge**: It allows you to pass the clean markdown content from parsing directly to extraction\n",
"2. **Markdown improves extraction**: Pre-processed markdown provides much better context than raw PDFs\n",
"3. **Classification enables smart routing**: Different document types can use different extraction schemas\n",
"4. **End-to-end automation**: The entire workflow can be automated for production use\n",
"\n",
"This approach is ideal for production document processing pipelines where you need to:\n",
"- Process various document types automatically\n",
"- Extract structured data consistently\n",
"- Maintain high accuracy and reliability\n",
"- Handle documents at scale\n",
"\n",
"The combination of these three services provides a powerful, flexible document processing pipeline that can handle complex, real-world document processing requirements."
]
}
],
"metadata": {
"kernelspec": {
"display_name": "llama_parse",
"display_name": "Python 3 (ipykernel)",
"language": "python",
"name": "llama_parse"
"name": "python3"
},
"language_info": {
"codemirror_mode": {
+1 -3
View File
@@ -1,6 +1,5 @@
from llama_cloud_services.parse import LlamaParse
from llama_cloud_services.extract import LlamaExtract, ExtractionAgent
from llama_cloud_services.utils import SourceText, FileInput
from llama_cloud_services.extract import LlamaExtract, ExtractionAgent, SourceText
from llama_cloud_services.constants import EU_BASE_URL
from llama_cloud_services.index import (
LlamaCloudCompositeRetriever,
@@ -13,7 +12,6 @@ __all__ = [
"LlamaExtract",
"ExtractionAgent",
"SourceText",
"FileInput",
"EU_BASE_URL",
"LlamaCloudIndex",
"LlamaCloudRetriever",
@@ -1,10 +0,0 @@
from llama_cloud_services.beta.classifier.client import ClassifyClient
from llama_cloud_services.beta.classifier.types import ClassifyJobResultsWithFiles
from llama_cloud_services.utils import SourceText, FileInput
__all__ = [
"ClassifyClient",
"ClassifyJobResultsWithFiles",
"SourceText",
"FileInput",
]
+27 -145
View File
@@ -1,7 +1,6 @@
import asyncio
import time
import warnings
from typing import Optional, List, Union
from typing import Optional
from pydantic import BaseModel
from llama_cloud.client import AsyncLlamaCloud
from llama_cloud.types import (
@@ -15,11 +14,7 @@ from llama_cloud.types import (
from llama_cloud.resources.classifier.client import OMIT
from llama_cloud_services.files.client import FileClient
from llama_cloud_services.constants import POLLING_TIMEOUT_SECONDS
from llama_cloud_services.utils import (
is_terminal_status,
augment_async_errors,
FileInput,
)
from llama_cloud_services.utils import is_terminal_status, augment_async_errors
from llama_index.core.async_utils import DEFAULT_NUM_WORKERS, run_jobs
from llama_cloud_services.beta.classifier.types import (
ClassifyJobResultsWithFiles,
@@ -171,98 +166,6 @@ class ClassifyClient:
)
)
async def aclassify(
self,
rules: list[ClassifierRule],
files: Union[FileInput, List[FileInput]],
parsing_configuration: Optional[ClassifyParsingConfiguration] = None,
raise_on_error: bool = True,
workers: int = DEFAULT_NUM_WORKERS,
show_progress: bool = False,
) -> ClassifyJobResultsWithFiles:
"""
Classify one or more files from various input types.
Args:
rules: The rules to use for classification.
files: The file(s) to classify. Can be a single file or list of files. Each can be:
- str/Path: File path
- SourceText: Text content or file with explicit filename
- File: Already uploaded file
- BufferedIOBase: File-like object
parsing_configuration: The parsing configuration to use for classification.
raise_on_error: Whether to raise an error if the classification job fails.
workers: Number of parallel workers for uploading files.
show_progress: Whether to show progress bars.
Returns:
The results of the classification job with file metadata.
"""
# Normalize to list
if not isinstance(files, list):
files = [files]
# Upload all files
coroutines = [
self.file_client.upload_content(file_input) for file_input in files
]
uploaded_files: List[File] = await run_jobs(
coroutines,
show_progress=show_progress,
workers=workers,
desc="Uploading files for classification",
)
# Classify
results = await self.aclassify_file_ids(
rules,
[file.id for file in uploaded_files],
parsing_configuration,
raise_on_error,
)
return ClassifyJobResultsWithFiles.from_classify_job_results(
results, uploaded_files
)
def classify(
self,
rules: list[ClassifierRule],
files: Union[FileInput, List[FileInput]],
parsing_configuration: Optional[ClassifyParsingConfiguration] = None,
raise_on_error: bool = True,
workers: int = DEFAULT_NUM_WORKERS,
show_progress: bool = False,
) -> ClassifyJobResultsWithFiles:
"""
Classify one or more files from various input types (synchronous version).
Args:
rules: The rules to use for classification.
files: The file(s) to classify. Can be a single file or list of files. Each can be:
- str/Path: File path
- SourceText: Text content or file with explicit filename
- File: Already uploaded file
- BufferedIOBase: File-like object
parsing_configuration: The parsing configuration to use for classification.
raise_on_error: Whether to raise an error if the classification job fails.
workers: Number of parallel workers for uploading files.
show_progress: Whether to show progress bars.
Returns:
The results of the classification job with file metadata.
"""
with augment_async_errors():
return asyncio.run(
self.aclassify(
rules,
files,
parsing_configuration,
raise_on_error,
workers,
show_progress,
)
)
async def aclassify_file_path(
self,
rules: list[ClassifierRule],
@@ -270,17 +173,11 @@ class ClassifyClient:
parsing_configuration: Optional[ClassifyParsingConfiguration] = None,
raise_on_error: bool = True,
) -> ClassifyJobResultsWithFiles:
"""
Deprecated: Use aclassify() instead.
"""
warnings.warn(
"aclassify_file_path is deprecated, use aclassify() instead",
DeprecationWarning,
stacklevel=2,
)
return await self.aclassify(
rules, file_input_path, parsing_configuration, raise_on_error
file = await self.file_client.upload_file(file_input_path)
results = await self.aclassify_file_ids(
rules, [file.id], parsing_configuration, raise_on_error
)
return ClassifyJobResultsWithFiles.from_classify_job_results(results, [file])
def classify_file_path(
self,
@@ -289,17 +186,12 @@ class ClassifyClient:
parsing_configuration: Optional[ClassifyParsingConfiguration] = None,
raise_on_error: bool = True,
) -> ClassifyJobResultsWithFiles:
"""
Deprecated: Use classify() instead.
"""
warnings.warn(
"classify_file_path is deprecated, use classify() instead",
DeprecationWarning,
stacklevel=2,
)
return self.classify(
rules, file_input_path, parsing_configuration, raise_on_error
)
with augment_async_errors():
return asyncio.run(
self.aclassify_file_path(
rules, file_input_path, parsing_configuration, raise_on_error
)
)
async def aclassify_file_paths(
self,
@@ -310,22 +202,17 @@ class ClassifyClient:
workers: int = DEFAULT_NUM_WORKERS,
show_progress: bool = False,
) -> ClassifyJobResultsWithFiles:
"""
Deprecated: Use aclassify() instead.
"""
warnings.warn(
"aclassify_file_paths is deprecated, use aclassify() instead",
DeprecationWarning,
stacklevel=2,
coroutines = [self.file_client.upload_file(path) for path in file_input_paths]
files: list[File] = await run_jobs(
coroutines,
show_progress=show_progress,
workers=workers,
desc="Uploading files for classification",
)
return await self.aclassify(
rules,
file_input_paths,
parsing_configuration,
raise_on_error,
workers,
show_progress,
results = await self.aclassify_file_ids(
rules, [file.id for file in files], parsing_configuration, raise_on_error
)
return ClassifyJobResultsWithFiles.from_classify_job_results(results, files)
def classify_file_paths(
self,
@@ -334,17 +221,12 @@ class ClassifyClient:
parsing_configuration: Optional[ClassifyParsingConfiguration] = None,
raise_on_error: bool = True,
) -> ClassifyJobResultsWithFiles:
"""
Deprecated: Use classify() instead.
"""
warnings.warn(
"classify_file_paths is deprecated, use classify() instead",
DeprecationWarning,
stacklevel=2,
)
return self.classify(
rules, file_input_paths, parsing_configuration, raise_on_error
)
with augment_async_errors():
return asyncio.run(
self.aclassify_file_paths(
rules, file_input_paths, parsing_configuration, raise_on_error
)
)
async def wait_for_job_completion(self, job_id: str) -> ClassifyJob:
"""
+1 -2
View File
@@ -2,16 +2,15 @@ from llama_cloud_services.extract.extract import (
LlamaExtract,
ExtractConfig,
ExtractionAgent,
SourceText,
ExtractTarget,
ExtractMode,
)
from llama_cloud_services.utils import SourceText, FileInput
__all__ = [
"LlamaExtract",
"ExtractionAgent",
"SourceText",
"FileInput",
"ExtractConfig",
"ExtractTarget",
"ExtractMode",
+100 -7
View File
@@ -2,9 +2,10 @@ import asyncio
import base64
import os
import time
from io import BufferedIOBase, TextIOWrapper
from io import BufferedIOBase, BufferedReader, BytesIO, TextIOWrapper
from pathlib import Path
from typing import List, Optional, Type, Union, Coroutine, Any, TypeVar
import secrets
import warnings
import httpx
from pydantic import BaseModel
@@ -32,8 +33,7 @@ from llama_cloud_services.extract.utils import (
JSONObjectType,
ExperimentalWarning,
)
from llama_cloud_services.utils import augment_async_errors, SourceText, FileInput
from llama_cloud_services.files.client import FileClient
from llama_cloud_services.utils import augment_async_errors
from llama_index.core.schema import BaseComponent
from llama_index.core.async_utils import run_jobs
from llama_index.core.bridge.pydantic import Field, PrivateAttr
@@ -188,6 +188,46 @@ async def _wait_for_job_result(
)
class SourceText:
def __init__(
self,
*,
file: Union[bytes, BufferedIOBase, TextIOWrapper, str, Path, None] = None,
text_content: Optional[str] = None,
filename: Optional[str] = None,
):
self.file = file
self.filename = filename
self.text_content = text_content
self._validate()
def _validate(self) -> None:
"""Ensure filename is provided when needed."""
if not ((self.file is None) ^ (self.text_content is None)):
raise ValueError("Either file or text_content must be provided.")
if self.text_content is not None:
if not self.filename:
random_hex = secrets.token_hex(4)
self.filename = f"text_input_{random_hex}.txt"
return
if isinstance(self.file, (bytes, BufferedIOBase, TextIOWrapper)):
if not self.filename and hasattr(self.file, "name"):
self.filename = os.path.basename(str(self.file.name))
elif not hasattr(self.file, "name") and self.filename is None:
raise ValueError(
"filename must be provided when file is bytes or a file-like object without a name"
)
elif isinstance(self.file, (str, Path)):
if not self.filename:
self.filename = os.path.basename(str(self.file))
else:
raise ValueError(f"Unsupported file type: {type(self.file)}")
FileInput = Union[str, Path, BufferedIOBase, SourceText, File]
def run_in_thread(
coro: Coroutine[Any, Any, T],
thread_pool: ThreadPoolExecutor,
@@ -280,7 +320,6 @@ class ExtractionAgent:
self._thread_pool = ThreadPoolExecutor(
max_workers=min(10, (os.cpu_count() or 1) + 4)
)
self._file_client = FileClient(client, project_id, organization_id)
@property
def id(self) -> str:
@@ -330,11 +369,65 @@ class ExtractionAgent:
ValueError: If filename is not provided for bytes input or for file-like objects
without a name attribute.
"""
return await self._file_client.upload_content(file_input)
file_contents: Optional[Union[BufferedIOBase, BytesIO]] = None
try:
if file_input.text_content is not None:
# Handle direct text content
file_contents = BytesIO(file_input.text_content.encode("utf-8"))
elif isinstance(file_input.file, TextIOWrapper):
# Handle text-based IO objects
file_contents = BytesIO(file_input.file.read().encode("utf-8"))
elif isinstance(file_input.file, (str, Path)):
# Handle file paths
file_contents = open(file_input.file, "rb")
elif isinstance(file_input.file, bytes):
# Handle bytes
file_contents = BytesIO(file_input.file)
elif isinstance(file_input.file, BufferedIOBase):
# Handle binary IO objects
file_contents = file_input.file
else:
raise ValueError(f"Unsupported file type: {type(file_input.file)}")
# Add name attribute to file object if needed
if not hasattr(file_contents, "name"):
file_contents.name = file_input.filename # type: ignore
return await self._client.files.upload_file(
project_id=self._project_id, upload_file=file_contents
)
finally:
if file_contents is not None and isinstance(
file_contents, (BufferedReader, BytesIO)
):
file_contents.close()
async def _upload_file(self, file_input: FileInput) -> File:
"""Upload a file from various input types using FileClient."""
return await self._file_client.upload_content(file_input)
source_text = None
if isinstance(file_input, File):
return file_input
if isinstance(file_input, SourceText):
source_text = file_input
elif isinstance(file_input, (str, Path)):
path = Path(file_input)
source_text = SourceText(file=path, filename=path.name)
else:
# Try to get filename from the file object if not provided
filename = None
if hasattr(file_input, "name"):
filename = os.path.basename(str(file_input.name))
if filename is None:
raise ValueError(
"Use SourceText to provide filename when uploading bytes or file-like objects."
)
warnings.warn(
"Use SourceText instead of bytes or file-like objects",
DeprecationWarning,
)
source_text = SourceText(file=file_input, filename=filename)
return await self.upload_file(source_text)
async def _wait_for_job_result(self, job_id: str) -> Optional[ExtractRun]:
"""Wait for and return the results of an extraction job."""
-82
View File
@@ -1,11 +1,9 @@
from io import BytesIO
from typing import BinaryIO
import os
from pathlib import Path
from llama_cloud.client import AsyncLlamaCloud
from llama_cloud.types import File, FileCreate
from typing import Optional
from llama_cloud_services.utils import SourceText, FileInput
class FileClient:
@@ -97,83 +95,3 @@ class FileClient:
project_id=self.project_id,
organization_id=self.organization_id,
)
async def upload_content(
self, file_input: FileInput, external_file_id: Optional[str] = None
) -> File:
"""
Upload content from various input types or fetch an already-uploaded file.
Args:
file_input: The content to upload. Can be:
- File: Already uploaded file (returned as-is)
- str/Path: Path to a file on disk
- SourceText: Text content, file, or file_id with explicit filename
- BufferedIOBase: File-like binary object
external_file_id: Optional external identifier for the file
Returns:
File: The uploaded (or fetched) file object
Raises:
ValueError: If the input type is not supported or required info is missing
"""
# If already a File object, return it
if isinstance(file_input, File):
return file_input
# Handle SourceText
if isinstance(file_input, SourceText):
# If file_id is provided, fetch the file object
if file_input.file_id is not None:
return await self.get_file(file_input.file_id)
elif file_input.text_content is not None:
# Handle direct text content
text_bytes = file_input.text_content.encode("utf-8")
return await self.upload_bytes(
text_bytes, external_file_id or file_input.filename or "file"
)
elif isinstance(file_input.file, (str, Path)):
# Handle file paths using the existing upload_file method
return await self.upload_file(
str(file_input.file), external_file_id or file_input.filename
)
elif isinstance(file_input.file, bytes):
# Handle bytes
return await self.upload_bytes(
file_input.file, external_file_id or file_input.filename or "file"
)
elif hasattr(file_input.file, "read"):
# Handle any file-like object (TextIOWrapper, BytesIO, BufferedReader, BufferedIOBase, etc.)
content = file_input.file.read() # type: ignore
if isinstance(content, str):
content = content.encode("utf-8")
return await self.upload_bytes(
content, external_file_id or file_input.filename or "file"
)
else:
raise ValueError(f"Unsupported file type: {type(file_input.file)}")
# Handle string/Path directly
elif isinstance(file_input, (str, Path)):
return await self.upload_file(str(file_input), external_file_id)
# Handle raw file-like objects
elif hasattr(file_input, "read"):
if hasattr(file_input, "name"):
filename = os.path.basename(str(file_input.name))
else:
filename = external_file_id or "file"
# Read content to determine size
content = file_input.read()
if isinstance(content, str):
content = content.encode("utf-8")
return await self.upload_bytes(content, external_file_id or filename)
else:
raise ValueError(
f"Unsupported file input type: {type(file_input)}. "
f"Supported types: str, Path, SourceText, BufferedIOBase, or File."
)
+2 -102
View File
@@ -3,14 +3,11 @@ import importlib.metadata
from contextlib import contextmanager
from typing import Generator
import difflib
from llama_cloud.types import StatusEnum, File
from llama_cloud.types import StatusEnum
import httpx
import packaging.version
from pydantic import BaseModel
from typing import Any, Dict, List, Tuple, Type, Union, Optional
from io import BufferedIOBase, TextIOWrapper
from pathlib import Path
import secrets
from typing import Any, Dict, List, Tuple, Type
# Asyncio error messages
nest_asyncio_err = "cannot be called from a running event loop"
@@ -107,100 +104,3 @@ def augment_async_errors() -> Generator[None, None, None]:
if nest_asyncio_err in str(e):
raise RuntimeError(nest_asyncio_msg)
raise
class SourceText:
"""
A wrapper class for providing text or file input with optional filename specification.
This class allows you to provide input in multiple ways:
- Direct text content via text_content parameter
- File paths as strings or Path objects
- Raw bytes
- File-like objects (BufferedIOBase, TextIOWrapper)
- Already-uploaded file ID via file_id parameter
Args:
file: The file input (bytes, file-like object, str path, or Path).
Mutually exclusive with text_content and file_id.
text_content: Raw text content to process. Mutually exclusive with file and file_id.
file_id: ID of an already-uploaded file. Mutually exclusive with file and text_content.
filename: Optional filename. Required for bytes/file-like objects without names.
If not provided, will be auto-generated for text_content or inferred from paths.
Examples:
# Direct text input
source = SourceText(text_content="Hello world")
# File path
source = SourceText(file="document.pdf")
# Bytes with filename
source = SourceText(file=b"...", filename="document.pdf")
# File-like object (will read from current position)
with open("document.pdf", "rb") as f:
source = SourceText(file=f)
# Already-uploaded file
source = SourceText(file_id="file_abc123")
"""
def __init__(
self,
*,
file: Union[bytes, BufferedIOBase, TextIOWrapper, str, Path, None] = None,
text_content: Optional[str] = None,
file_id: Optional[str] = None,
filename: Optional[str] = None,
):
self.file = file
self.filename = filename
self.text_content = text_content
self.file_id = file_id
self._validate()
def _validate(self) -> None:
"""Ensure filename is provided when needed."""
# Check that exactly one of file, text_content, or file_id is provided
provided = sum(
[
self.file is not None,
self.text_content is not None,
self.file_id is not None,
]
)
if provided == 0:
raise ValueError("One of file, text_content, or file_id must be provided.")
elif provided > 1:
raise ValueError(
"Only one of file, text_content, or file_id can be provided."
)
# If file_id is provided, we don't need filename validation
if self.file_id is not None:
return
if self.text_content is not None:
if not self.filename:
random_hex = secrets.token_hex(4)
self.filename = f"text_input_{random_hex}.txt"
return
if isinstance(self.file, (bytes, BufferedIOBase, TextIOWrapper)):
if not self.filename and hasattr(self.file, "name"):
self.filename = os.path.basename(str(self.file.name))
elif self.filename is None and not hasattr(self.file, "name"):
raise ValueError(
"filename must be provided when file is bytes or a file-like object without a name"
)
elif isinstance(self.file, (str, Path)):
if not self.filename:
self.filename = os.path.basename(str(self.file))
else:
raise ValueError(f"Unsupported file type: {type(self.file)}")
# Type alias for file input that can be used across services
FileInput = Union[str, Path, BufferedIOBase, SourceText, File]
Generated
+2 -2
View File
@@ -1,5 +1,5 @@
version = 1
revision = 2
revision = 3
requires-python = ">=3.9, <4.0"
resolution-markers = [
"python_full_version >= '3.14'",
@@ -1596,7 +1596,7 @@ wheels = [
[[package]]
name = "llama-cloud-services"
version = "0.6.76"
version = "0.6.73"
source = { editable = "." }
dependencies = [
{ name = "click", version = "8.1.8", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version < '3.10'" },
-6
View File
@@ -1,11 +1,5 @@
# llama-cloud-services
## 0.3.10
### Patch Changes
- fee516d: Adding LlamaClassify among the available LlamaCloud services
## 0.3.9
### Patch Changes
@@ -1,8 +0,0 @@
{
"type": "module",
"main": "./dist/index.cjs",
"module": "./dist/index.js",
"types": "./dist/index.d.ts",
"exports": "./dist/index.js",
"private": true
}
+2 -14
View File
@@ -1,6 +1,6 @@
{
"name": "llama-cloud-services",
"version": "0.3.10",
"version": "0.3.9",
"type": "module",
"license": "MIT",
"scripts": {
@@ -24,8 +24,7 @@
"./reader",
"./parse",
"./beta/agent",
"./extract",
"./classify"
"./extract"
],
"exports": {
"./openapi.json": "./openapi.json",
@@ -84,17 +83,6 @@
},
"default": "./extract/dist/index.js"
},
"./classify": {
"require": {
"types": "./classify/dist/index.d.cts",
"default": "./classify/dist/index.cjs"
},
"import": {
"types": "./classify/dist/index.d.ts",
"default": "./classify/dist/index.js"
},
"default": "./classify/dist/index.js"
},
".": {
"require": {
"types": "./dist/index.d.cts",
@@ -1,69 +0,0 @@
import { createClient, createConfig, type Client } from "@hey-api/client-fetch";
import {
classify,
type ClassifyParsingConfiguration,
type ClassifierRule,
type ClassifyJobResults,
} from "./classify";
import { getUrl } from "./utils";
import { getEnv } from "@llamaindex/env";
import { File } from "buffer";
export class LlamaClassify {
private client: Client;
constructor(
apiKey: string | undefined = undefined,
baseUrl: string | undefined = undefined,
region: string | undefined = undefined,
) {
const key = apiKey ?? getEnv("LLAMA_CLOUD_API_KEY");
if (typeof key === "undefined") {
throw new Error(
"No API key provided and no API key found in environment. Please pass the API key or set `LLAMA_CLOUD_API_KEY` as an environment variable.",
);
}
const url = getUrl(baseUrl, region);
this.client = createClient(
createConfig({
baseUrl: url,
headers: {
Authorization: `Bearer ${key}`,
},
}),
);
}
async classify(
rules: ClassifierRule[],
parsingConfiguration: ClassifyParsingConfiguration,
fileContents:
| Buffer<ArrayBufferLike>[]
| File[]
| Uint8Array<ArrayBuffer>[]
| string[]
| undefined = undefined,
filePaths: string[] | undefined = undefined,
projectId: string | null = null,
organizationId: string | null = null,
pollingInterval: number = 1,
maxPollingIterations: number = 1800,
maxRetriesOnError: number = 10,
retryInterval: number = 0.5,
): Promise<ClassifyJobResults> {
const result = await classify(
rules,
parsingConfiguration,
fileContents,
filePaths,
projectId,
organizationId,
this.client,
pollingInterval,
maxPollingIterations,
maxRetriesOnError,
retryInterval,
);
return result;
}
}
+19 -1
View File
@@ -4,7 +4,25 @@ import * as extract from "./extract";
import type { ExtractAgent, ExtractConfig } from "./extract";
import { getEnv } from "@llamaindex/env";
import type { ExtractResult } from "./type";
import { getUrl } from "./utils";
const URLS = {
us: "https://api.cloud.llamaindex.ai",
eu: "https://api.cloud.eu.llamaindex.ai",
"us-staging": "https://api.staging.llamaindex.ai",
} as const;
function getUrl(baseUrl: string | undefined, region: string | undefined) {
if (typeof baseUrl != "undefined") {
return baseUrl;
}
if (typeof region === "undefined") {
return URLS["us"];
} else if (region === "us" || region === "eu" || region === "us-staging") {
return URLS[region];
} else {
throw new Error(`Unsupported region: ${region}`);
}
}
export class LlamaExtractAgent {
private agent: ExtractAgent;
-289
View File
@@ -1,289 +0,0 @@
import type {
Options,
CreateClassifyJobApiV1ClassifierJobsPostData,
ClassifyJobCreate,
ClassifierRule,
ClassifyParsingConfiguration,
GetClassifyJobApiV1ClassifierJobsClassifyJobIdGetData,
GetClassificationJobResultsApiV1ClassifierJobsClassifyJobIdResultsGetData,
ClassifyJobResults,
} from "./api";
import {
StatusEnum,
createClassifyJobApiV1ClassifierJobsPost,
getClassifyJobApiV1ClassifierJobsClassifyJobIdGet,
getClassificationJobResultsApiV1ClassifierJobsClassifyJobIdResultsGet,
} from "./api";
import type { Client } from "@hey-api/client-fetch";
import { sleep } from "./utils";
import { uploadFile } from "./fileUpload";
import { File } from "buffer";
async function createClassifyJob(
fileIds: string[],
rules: ClassifierRule[],
parsingConfiguration: ClassifyParsingConfiguration,
organizationId: null | string,
projectId: null | string,
client: Client | undefined,
maxRetriesOnError: number = 10,
retryInterval: number = 0.5,
): Promise<string> {
const rawData = {
file_ids: fileIds,
rules: rules,
parsing_configuration: parsingConfiguration,
} as ClassifyJobCreate;
const data = {
body: rawData,
query: {
project_id: projectId,
organization_id: organizationId,
},
} as CreateClassifyJobApiV1ClassifierJobsPostData;
const options = data as Options<CreateClassifyJobApiV1ClassifierJobsPostData>;
if (typeof client != "undefined") {
options.client = client;
}
let retries = 0;
while (true) {
if (retries > maxRetriesOnError) {
throw new Error(
"Error while creating the classify job: Exceeded maximum number of retries, the API keeps returning errors.",
);
}
const response = await createClassifyJobApiV1ClassifierJobsPost(options);
if (!response.response.ok) {
if ("error" in response) {
console.log(
`An error occurred while creating the classification job.\nDetails:\n\n${JSON.stringify(
response.error,
)}\n\nRetrying...`,
);
}
retries++;
await sleep(retryInterval * 1000);
} else {
if (typeof response.data != "undefined") {
return response.data.id;
} else {
throw new Error(
"Error while creating the classify job: the job creation succeeded but no data where returned",
);
}
}
}
}
async function pollForJobCompletion(
jobId: string,
interval: number = 1,
maxIterations: number = 1800,
client: Client | undefined = undefined,
): Promise<boolean> {
let status: StatusEnum | undefined = undefined;
const jobData = {
path: { classify_job_id: jobId },
} as GetClassifyJobApiV1ClassifierJobsClassifyJobIdGetData;
const jobOptions =
jobData as Options<GetClassifyJobApiV1ClassifierJobsClassifyJobIdGetData>;
if (typeof client != "undefined") {
jobOptions.client = client;
}
let numIterations: number = 0;
while (true) {
if (numIterations > maxIterations) {
return false;
}
const response =
await getClassifyJobApiV1ClassifierJobsClassifyJobIdGet(jobOptions);
if (!response.response.ok) {
numIterations++;
}
if (typeof response.data != "undefined") {
status = response.data.status as StatusEnum;
if (status == StatusEnum.CANCELLED || status == StatusEnum.ERROR) {
throw new Error("There was an error during the classification job.");
} else if (status == StatusEnum.SUCCESS) {
return true;
} else {
numIterations++;
await sleep(interval * 1000);
}
}
}
}
async function getJobResult(
jobId: string,
client: Client | undefined = undefined,
projectId: string | null = null,
organizationId: string | null = null,
maxRetriesOnError: number = 10,
retryInterval: number = 0.5,
): Promise<ClassifyJobResults> {
const jobData = {
path: { classify_job_id: jobId },
query: { organization_id: organizationId, project_id: projectId },
} as GetClassificationJobResultsApiV1ClassifierJobsClassifyJobIdResultsGetData;
const jobOptions =
jobData as Options<GetClassificationJobResultsApiV1ClassifierJobsClassifyJobIdResultsGetData>;
if (typeof client != "undefined") {
jobOptions.client = client;
}
let retries: number = 0;
while (true) {
if (retries > maxRetriesOnError) {
throw new Error(
"Error while getting the result of the classification job: Exceeded maximum number of retries, the API keeps returning errors.",
);
}
const response =
await getClassificationJobResultsApiV1ClassifierJobsClassifyJobIdResultsGet(
jobOptions,
);
if (!response.response.ok) {
if ("error" in response) {
console.log(
"An error occurred: ",
JSON.stringify(response.error),
"\nRetrying...",
);
}
retries++;
await sleep(retryInterval * 1000);
}
if (typeof response.data != "undefined") {
return response.data as ClassifyJobResults;
} else {
throw new Error(
"Error while retrieving results for the classify job: the result was successfully obtained but no data were returned",
);
}
}
}
export async function classify(
rules: ClassifierRule[],
parsingConfiguration: ClassifyParsingConfiguration,
fileContents:
| Buffer<ArrayBufferLike>[]
| File[]
| Uint8Array<ArrayBuffer>[]
| string[]
| undefined = undefined,
filePaths: string[] | undefined = undefined,
projectId: string | null = null,
organizationId: string | null = null,
client: Client | undefined = undefined,
pollingInterval: number = 1,
maxPollingIterations: number = 1800,
maxRetriesOnError: number = 10,
retryInterval: number = 0.5,
): Promise<ClassifyJobResults> {
const fileIds: string[] = [];
if (!filePaths && !fileContents) {
throw new Error(
"One between filePath and fileContent needs to be provided",
);
}
if (filePaths) {
const uploadPromises = filePaths.map(async (name) => {
try {
const fileId = await uploadFile(
name,
undefined,
undefined,
projectId,
organizationId,
client,
maxRetriesOnError,
retryInterval,
);
if (fileId) {
return fileId;
} else {
console.error(`Unable to upload ${name}, skipping...`);
return null;
}
} catch (error) {
console.error(`Error uploading ${name}:`, error);
return null;
}
});
const results = await Promise.all(uploadPromises);
fileIds.push(...results.filter((id) => id !== null));
}
if (fileContents) {
const uploadPromises = fileContents.map(async (content) => {
try {
const fileId = await uploadFile(
undefined,
content,
undefined,
projectId,
organizationId,
client,
maxRetriesOnError,
retryInterval,
);
if (fileId) {
return fileId;
} else {
console.error(`Unable to upload file (content), skipping...`);
return null;
}
} catch (error) {
console.error(`Error uploading file (content):`, error);
return null;
}
});
const results = await Promise.all(uploadPromises);
fileIds.push(...results.filter((id) => id !== null));
}
if (fileIds.length == 0) {
throw new Error(
"None of the provided files was successfully uploaded, it is not possible to create a classification job.",
);
}
const jobId = await createClassifyJob(
fileIds,
rules,
parsingConfiguration,
organizationId,
projectId,
client,
maxRetriesOnError,
retryInterval,
);
const success = await pollForJobCompletion(
jobId,
pollingInterval,
maxPollingIterations,
client,
);
if (!success) {
throw new Error("Your job is taking longer than 10 minutes, timing out...");
} else {
return (await getJobResult(
jobId,
client,
projectId,
organizationId,
maxRetriesOnError,
retryInterval,
)) as ClassifyJobResults;
}
}
export {
type ClassifierRule,
type ClassifyJobResults,
type ClassifyParsingConfiguration,
};
+100 -1
View File
@@ -1,5 +1,9 @@
import { emitWarning } from "process";
import fs from "fs/promises";
import { Blob } from "buffer";
import * as path from "path";
import type { ExtractResult } from "./type";
import { randomUUID } from "@llamaindex/env";
import { File } from "buffer";
import {
type Options,
@@ -15,6 +19,7 @@ import {
type GetJobApiV1ExtractionJobsJobIdGetData,
type GetJobResultApiV1ExtractionJobsJobIdResultGetData,
StatusEnum,
type UploadFileApiV1FilesPostData,
type StatelessExtractionRequest,
type ExtractStatelessApiV1ExtractionRunPostData,
type DeleteExtractionAgentApiV1ExtractionExtractionAgentsExtractionAgentIdDeleteData,
@@ -24,12 +29,17 @@ import {
runJobApiV1ExtractionJobsPost,
getJobApiV1ExtractionJobsJobIdGet,
getJobResultApiV1ExtractionJobsJobIdResultGet,
uploadFileApiV1FilesPost,
extractStatelessApiV1ExtractionRunPost,
deleteExtractionAgentApiV1ExtractionExtractionAgentsExtractionAgentIdDelete,
} from "./api";
import type { Client } from "@hey-api/client-fetch";
import { sleep } from "./utils";
import { uploadFile } from "./fileUpload";
import { fileTypeFromBuffer } from "file-type";
type BodyUploadFileApiV1FilesPost = {
upload_file: Blob | File;
};
export async function createAgent(
name: string,
@@ -211,6 +221,95 @@ export async function getAgent(
}
}
function textToFile(text: string, fileName: string | null = null) {
return new File(
[text],
fileName ?? "uploadedFile_" + randomUUID().replaceAll("-", "_") + ".txt",
);
}
async function uploadFile(
filePath: string | undefined = undefined,
fileContent:
| Buffer<ArrayBufferLike>
| File
| Uint8Array<ArrayBuffer>
| string
| undefined = undefined,
fileName: string | undefined = undefined,
project_id: string | null = null,
organization_id: string | null = null,
client: Client | undefined = undefined,
maxRetriesOnError: number = 10,
retryInterval: number = 0.5,
): Promise<string | undefined> {
let file: File | undefined = undefined;
if (typeof filePath === "undefined" && typeof fileContent === "undefined") {
throw new Error(
"One between filePath and fileContent needs to be provided",
);
} else if (typeof filePath != "undefined") {
const buffer = await fs.readFile(filePath);
const actualFileName = fileName ?? path.basename(filePath);
const uint8Array = new Uint8Array(buffer);
file = new File([uint8Array], actualFileName);
} else if (typeof fileContent != "undefined") {
if (fileContent instanceof File) {
file = fileContent;
} else if (fileContent instanceof Buffer) {
const fileType = await fileTypeFromBuffer(fileContent);
const ext = fileType?.ext ?? "pdf";
const uint8Array = new Uint8Array(fileContent);
file = new File(
[uint8Array],
fileName ??
"uploadedFile_" + randomUUID().replaceAll("-", "_") + "." + ext,
);
} else if (fileContent instanceof Uint8Array) {
const fileType = await fileTypeFromBuffer(fileContent);
const ext = fileType?.ext ?? "pdf";
file = new File(
[fileContent],
fileName ??
"uploadedFile_" + randomUUID().replaceAll("-", "_") + "." + ext,
);
} else if (typeof fileContent === "string") {
file = textToFile(fileContent, fileName);
} else {
throw new Error("Unsupported fileContent type");
}
}
const fileToUpload = {
upload_file: file,
} as BodyUploadFileApiV1FilesPost;
const uploadData = {
body: fileToUpload,
query: { organization_id: organization_id, project_id: project_id },
} as UploadFileApiV1FilesPostData;
const uploadOptions = uploadData as Options<UploadFileApiV1FilesPostData>;
if (typeof client != "undefined") {
uploadOptions.client = client;
}
let retries: number = 0;
while (true) {
if (retries > maxRetriesOnError) {
throw new Error(
"Error while processing your file: Exceeded maximum number of retries, the API keeps returning errors.",
);
}
const uploadResponse = await uploadFileApiV1FilesPost(uploadOptions);
let fileId: string | undefined = undefined;
if (!uploadResponse.response.ok) {
retries++;
await sleep(retryInterval * 1000);
}
if (typeof uploadResponse.data != "undefined") {
fileId = uploadResponse.data.id as string;
return fileId;
}
}
}
async function createExtractJob(
options:
| Options<RunJobApiV1ExtractionJobsPostData>
-109
View File
@@ -1,109 +0,0 @@
import fs from "fs/promises";
import { Blob } from "buffer";
import * as path from "path";
import { randomUUID } from "@llamaindex/env";
import { File } from "buffer";
import {
type Options,
type UploadFileApiV1FilesPostData,
uploadFileApiV1FilesPost,
} from "./api";
import type { Client } from "@hey-api/client-fetch";
import { sleep } from "./utils";
import { fileTypeFromBuffer } from "file-type";
type BodyUploadFileApiV1FilesPost = {
upload_file: Blob | File;
};
function textToFile(text: string, fileName: string | null = null) {
return new File(
[text],
fileName ?? "uploadedFile_" + randomUUID().replaceAll("-", "_") + ".txt",
);
}
export async function uploadFile(
filePath: string | undefined = undefined,
fileContent:
| Buffer<ArrayBufferLike>
| File
| Uint8Array<ArrayBuffer>
| string
| undefined = undefined,
fileName: string | undefined = undefined,
project_id: string | null = null,
organization_id: string | null = null,
client: Client | undefined = undefined,
maxRetriesOnError: number = 10,
retryInterval: number = 0.5,
): Promise<string | undefined> {
let file: File | undefined = undefined;
if (typeof filePath === "undefined" && typeof fileContent === "undefined") {
throw new Error(
"One between filePath and fileContent needs to be provided",
);
} else if (typeof filePath != "undefined") {
const buffer = await fs.readFile(filePath);
const actualFileName = fileName ?? path.basename(filePath);
const uint8Array = new Uint8Array(buffer);
file = new File([uint8Array], actualFileName);
} else if (typeof fileContent != "undefined") {
if (fileContent instanceof File) {
file = fileContent;
} else if (fileContent instanceof Buffer) {
const fileType = await fileTypeFromBuffer(fileContent);
const ext = fileType?.ext ?? "pdf";
const uint8Array = new Uint8Array(fileContent);
file = new File(
[uint8Array],
fileName ??
"uploadedFile_" + randomUUID().replaceAll("-", "_") + "." + ext,
);
} else if (fileContent instanceof Uint8Array) {
const fileType = await fileTypeFromBuffer(fileContent);
const ext = fileType?.ext ?? "pdf";
file = new File(
[fileContent],
fileName ??
"uploadedFile_" + randomUUID().replaceAll("-", "_") + "." + ext,
);
} else if (typeof fileContent === "string") {
file = textToFile(fileContent, fileName);
} else {
throw new Error("Unsupported fileContent type");
}
}
const fileToUpload = {
upload_file: file,
} as BodyUploadFileApiV1FilesPost;
const uploadData = {
body: fileToUpload,
query: { organization_id: organization_id, project_id: project_id },
} as UploadFileApiV1FilesPostData;
const uploadOptions = uploadData as Options<UploadFileApiV1FilesPostData>;
if (typeof client != "undefined") {
uploadOptions.client = client;
}
let retries: number = 0;
while (true) {
if (retries > maxRetriesOnError) {
throw new Error(
"Error while processing your file: Exceeded maximum number of retries, the API keeps returning errors.",
);
}
const uploadResponse = await uploadFileApiV1FilesPost(uploadOptions);
let fileId: string | undefined = undefined;
if (!uploadResponse.response.ok) {
retries++;
await sleep(retryInterval * 1000);
}
if (
uploadResponse.response.ok &&
typeof uploadResponse.data != "undefined"
) {
fileId = uploadResponse.data.id as string;
return fileId;
}
}
}
-6
View File
@@ -8,9 +8,3 @@ export type { CloudConstructorParams } from "./type.js";
export { LlamaParseReader } from "./reader.js";
export { LlamaExtract, LlamaExtractAgent } from "./LlamaExtract.js";
export type { ExtractConfig } from "./extract.js";
export { LlamaClassify } from "./LlamaClassify.js";
export type {
ClassifierRule,
ClassifyJobResults,
ClassifyParsingConfiguration,
} from "./classify.js";
-22
View File
@@ -117,25 +117,3 @@ export function getSavePath(downloadPath: string, i: number): string {
return savePath;
}
const URLS = {
us: "https://api.cloud.llamaindex.ai",
eu: "https://api.cloud.eu.llamaindex.ai",
"us-staging": "https://api.staging.llamaindex.ai",
} as const;
export function getUrl(
baseUrl: string | undefined,
region: string | undefined,
) {
if (typeof baseUrl != "undefined") {
return baseUrl;
}
if (typeof region === "undefined") {
return URLS["us"];
} else if (region === "us" || region === "eu" || region === "us-staging") {
return URLS[region];
} else {
throw new Error(`Unsupported region: ${region}`);
}
}
@@ -2,8 +2,6 @@ import { describe, it, expect, beforeEach, beforeAll } from "vitest";
import { LlamaParseReader } from "../src/reader.js";
import { LlamaCloudIndex } from "../src/LlamaCloudIndex.js";
import { LlamaExtract, LlamaExtractAgent } from "../src/LlamaExtract.js";
import { LlamaClassify } from "../src/LlamaClassify.js";
import { ClassifierRule, ClassifyParsingConfiguration } from "../src/classify.js";
import { Document } from "@llamaindex/core/schema";
import { fs } from "@llamaindex/env";
import { ExtractConfig } from "../src/api.js";
@@ -491,59 +489,6 @@ describe("Integration Tests", () => {
);
});
describe("LlamaClassify Integration", () => {
it.skipIf(skipIfNoApiKey)(
"should classify data correctly (file paths and file contents) ",
async () => {
const classifyClient = new LlamaClassify(
process.env.LLAMA_CLOUD_API_KEY!,
"https://api.cloud.llamaindex.ai",
);
const testContent =
`A Fox one day spied a beautiful bunch of ripe grapes hanging from a vine trained along the branches of a tree. The grapes seemed ready to burst with juice, and the Fox's mouth watered as he gazed longingly at them. The bunch hung from a high branch, and the Fox had to jump for it. The first time he jumped he missed it by a long way. So he walked off a short distance and took a running leap at it, only to fall short once more. Again and again he tried, but in vain. Now he sat down and looked at the grapes in disgust. "What a fool I am," he said. "Here I am wearing myself out to get a bunch of sour grapes that are not worth gaping for." And off he walked very, very scornfully.There are many who pretend to despise and belittle that which is beyond their reach.`;
const testFilePath = "the_fox_and_the_grapes.md";
await fs.writeFile(testFilePath, new TextEncoder().encode(testContent));
const rules: ClassifierRule[] = [
{type: "fable", description: "A short story featuring animals whose aim is to teach the reader a lesson (the moral of the story)"},
{type: "fairy_tale", description: "A mid-to-long story featuring humans, magic creatures and other characters, whose main aim is to entertain the readers."}
]
const parsingConfig: ClassifyParsingConfiguration = {lang: "en"}
const result = await classifyClient.classify(
rules,
parsingConfig,
undefined,
["the_fox_and_the_grapes.md"]
);
expect("items" in result).toBeTruthy();
expect(result.items.length).toBeGreaterThan(0);
expect("result" in result.items[0]).toBeTruthy();
expect(result.items[0].result!.type === "fable").toBeTruthy();
const buffer = await fs.readFile("the_fox_and_the_grapes.md");
const resultBuffer = await classifyClient.classify(
rules,
parsingConfig,
[buffer],
);
expect("items" in resultBuffer).toBeTruthy();
expect(resultBuffer.items.length).toBeGreaterThan(0);
expect("result" in resultBuffer.items[0]).toBeTruthy();
expect(resultBuffer.items[0].result!.type === "fable").toBeTruthy();
try {
await fs.unlink("the_fox_and_the_grapes.md")
} catch(err) {
console.log(`Unable to delete file the_fox_and_the_grapes.md because of ${err}`)
}
},
60000,
);
});
describe("LlamaExtract Integration", () => {
it.skipIf(skipIfNoApiKey)(
"should create agents correctly",