Ingests document metadata from cloud storage and collaboration platforms (S3, SharePoint, OneDrive, Google Drive, Box, Dropbox).
ingest_sourceIngestSourceOperator provides a unified interface for multiple document sources. It is a metadata-only operator — it discovers files and emits path/metadata rows; text content is extracted by a downstream ExtractOperator.
OperatorConstants.FileExtensions.BASE_EXTENSIONSIngest documents from one or more local directories or individual files. No credentials are required.
Configuration:
node_config = {
'provider': 'filesystem',
'connection_params': {
'paths': ['/data/invoices', '/data/contracts'],
'recursive': True,
'exclude_patterns': ['*.tmp', '__pycache__/*'],
'max_file_size_mb': 100,
'follow_symlinks': False
}
}
Parameters:
| Parameter | Type | Required | Default | Description |
|---|---|---|---|---|
paths |
list[str] |
Yes | — | One or more absolute or relative paths to files or directories |
recursive |
bool | No | True |
Recursively traverse subdirectories |
exclude_patterns |
list[str] | No | [] |
Glob patterns to skip (e.g. ["*.tmp", "__pycache__/*"]) |
max_file_size_mb |
int | No | None |
Skip files larger than this size (MB). None means no limit |
follow_symlinks |
bool | No | False |
Follow symbolic links during directory traversal |
File filtering is also controlled by the top-level include_filter / exclude_filter operator parameters (comma-separated extension list, e.g. "pdf,docx,txt").
Ingest documents from Amazon S3 buckets and S3-compatible storage services (IBM Cloud Object Storage, MinIO, etc.).
Configuration (AWS S3):
node_config = {
'provider': 's3',
'connection_params': {
'bucket': 'your-bucket-name',
'prefix': 'optional/path/prefix/', # Optional
'region': 'us-east-1' # Optional
},
'credentials': {
'access_key': 'YOUR_AWS_ACCESS_KEY',
'secret_key': 'YOUR_AWS_SECRET_KEY' # pragma: allowlist secret
}
}
Configuration (IBM Cloud Object Storage):
node_config = {
'provider': 's3',
'connection_params': {
'bucket': 'your-bucket-name',
'prefix': 'optional/path/prefix/', # Optional
'endpoint_url': 'https://s3.us-south.cloud-object-storage.appdomain.cloud'
},
'credentials': {
'access_key': 'YOUR_IBM_ACCESS_KEY',
'secret_key': 'YOUR_IBM_SECRET_KEY' # pragma: allowlist secret
}
}
Parameters:
bucket (required): S3 bucket nameprefix (optional): S3 key prefix to filter objects. Supports both directory-level and single file ingestion:
'documents/reports/' - ingests all files in the directory'documents/report.pdf' - ingests only the specified file'' (empty string) - ingests all files in the bucketendpoint_url (optional): Custom S3 endpoint URL for S3-compatible storage (e.g., IBM COS, MinIO). Leave empty for AWS S3.region (optional): AWS region (e.g., ‘us-east-1’). Optional for S3-compatible storage.access_key (required): AWS access key ID or S3-compatible access keysecret_key (required): AWS secret access key or S3-compatible secret keyrecursive (optional): Whether to recursively traverse subdirectories (default: True)exclude_patterns (optional): List of glob patterns to exclude (e.g., [‘*.tmp’, ‘.DS_Store’])max_file_size_mb (optional): Maximum file size in MB to processskip_hidden_files (optional): Whether to skip hidden files (default: True)skip_empty_files (optional): Whether to skip files with zero size (default: True)verify_expected_bucket_owner (optional): When True, verifies that the S3 bucket is owned by the caller’s AWS account via STS GetCallerIdentity. If the bucket owner does not match, AWS rejects the request. Default False. Has no effect for S3-compatible storage (IBM COS, MinIO).Note: File extension filtering is configured at the operator level using include_filter and exclude_filter parameters (see File Filtering section below).
Ingest documents from SharePoint document libraries.
Configuration:
node_config = {
'provider': 'sharepoint',
'connection_params': {
'document_library_id': 'your-library-id'
},
'credentials': {
'client_id': 'YOUR_CLIENT_ID',
'client_secret': 'YOUR_CLIENT_SECRET', # pragma: allowlist secret
'tenant_id': 'YOUR_TENANT_ID'
}
}
Prerequisites:
pip install O365Parameters:
document_library_id (required): SharePoint document library IDclient_id (required): Azure AD application client IDclient_secret (required): Azure AD application client secrettenant_id (required): Azure AD tenant IDIngest documents from OneDrive folders.
Configuration:
node_config = {
'provider': 'onedrive',
'connection_params': {
'drive_id': 'your-drive-id',
'folder_path': '/Documents/MyFolder' # Optional
},
'credentials': {
'client_id': 'YOUR_CLIENT_ID',
'client_secret': 'YOUR_CLIENT_SECRET', # pragma: allowlist secret
'tenant_id': 'YOUR_TENANT_ID'
}
}
Prerequisites:
pip install O365Parameters:
drive_id (required): OneDrive drive IDfolder_path (optional): Path to specific folderclient_id (required): Azure AD application client IDclient_secret (required): Azure AD application client secrettenant_id (required): Azure AD tenant IDIngest documents from Google Drive folders using OAuth 2.0 authentication.
Configuration:
node_config = {
'provider': 'google_drive',
'connection_params': {
'folder_id': 'your-folder-id',
'recursive': False # Optional: include subfolders
},
'credentials': {
'credentials_json_path': '/path/to/client_secret.json',
'token_path': '/path/to/token.json', # Optional
'scopes': ['https://www.googleapis.com/auth/drive.readonly'] # Optional
}
}
Prerequisites:
pip install google-auth-oauthlib google-auth-httplib2 google-api-python-clientParameters:
folder_id (required): Google Drive folder IDrecursive (optional): Boolean, include subfolders (default: False)credentials_json_path (required): Path to OAuth client secret JSONtoken_path (optional): Path to store OAuth tokens (default: ~/.credentials/token.json)scopes (optional): List of OAuth scopes (default: ['https://www.googleapis.com/auth/drive.readonly'])OAuth Scopes: The operator uses read-only access by default for security. Available scopes:
https://www.googleapis.com/auth/drive.readonly - Read-only access (recommended)https://www.googleapis.com/auth/drive - Full access to all fileshttps://www.googleapis.com/auth/drive.file - Per-file accessImportant: If you change scopes, you must delete the existing token file to re-authenticate with the new permissions.
Ingest documents from Box folders using JWT authentication.
Configuration:
node_config = {
'provider': 'box_driver',
'connection_params': {
'folder_id': '0', # Optional: Box folder ID to start from (default: '0' for root)
'recursive': True, # Optional: include subfolders
'max_file_size_mb': 50, # Optional: max file size in MB
'exclude_patterns': ['*.tmp', 'Trash/*'] # Optional: patterns to exclude
},
'credentials': {
'credentials_json_path': '/path/to/box_jwt_config.json'
},
'include_filter': 'pdf,docx,txt,pptx,xlsx', # Optional: file extensions to include
'max_files': 100 # Optional
}
Prerequisites:
pip install box-sdk-genParameters:
folder_id (optional): Box folder ID to start ingestion from (default: ‘0’ for root folder)recursive (optional): Boolean, include subfolders (default: False)max_file_size_mb (optional): Maximum file size in MB to processexclude_patterns (optional): List of glob patterns to exclude (e.g., ['*.tmp', 'Trash/*'])credentials_json_path (required): Path to Box JWT configuration JSON fileBox JWT Setup:
Authentication Flow: The adapter uses JWT (JSON Web Token) authentication which provides:
Security Notes:
Ingest documents from a Dropbox account using the Dropbox SDK.
Configuration:
node_config = {
'provider': 'dropbox',
'connection_params': {
'folder_path': '/Reports', # Optional: folder to ingest from ('' or '/' for account root)
'recursive': True, # Optional: include subfolders (default: True)
'max_file_size_mb': 50, # Optional: max file size in MB
'exclude_patterns': ['*.tmp', '*/Archive/*'] # Optional: patterns to exclude
},
'credentials': {
'access_token': '${DROPBOX_ACCESS_TOKEN}'
},
'include_filter': 'pdf,docx,txt', # Optional: file extensions to include
'max_files': 100 # Optional
}
Prerequisites:
account_info.read, files.metadata.read, files.content.readpip install dropboxParameters:
folder_path (optional): Dropbox folder to ingest from; empty string or / means the account rootfile_path (optional): Ingest a single file by path (/Reports/q1.pdf) or file id (id:abc123); ignores folder settings and filtersrecursive (optional): Boolean, include subfolders (default: True)max_file_size_mb (optional): Maximum file size in MB to processexclude_patterns (optional): List of glob patterns to exclude (e.g., ['*.tmp', '*/Archive/*'])access_token (credentials): Dropbox OAuth2 access tokenrefresh_token, app_key, app_secret (credentials): Long-lived alternative to access_tokenDropbox App Setup:
account_info.read, files.metadata.read and files.content.read, then submittoken_access_type=offline to obtain a refresh token${DROPBOX_ACCESS_TOKEN} in the flowAuthentication Notes:
Security Notes:
Extend functionality with custom LangChain-compatible loaders.
Configuration:
node_config = {
'provider': 'custom',
'connection_params': {
'loader_class_path': 'my_package.loaders.CustomLoader',
# Additional parameters specific to your loader
'param1': 'value1',
'param2': 'value2'
},
'credentials': {
# Credentials specific to your loader
'api_key': 'YOUR_API_KEY' # pragma: allowlist secret
}
}
Parameters:
loader_class_path (required): Python import path to loader class (e.g., my_package.loaders.FileNetLoader)__init__ methodRequirements:
load() methodload() method must return List[Document]from docpipe.core.operators.ingest.ingest_source import IngestSourceOperator
import pyarrow as pa
# Configure the operator
node_config = {
'provider': 's3',
'connection_params': {
'bucket': 'my-bucket',
'prefix': 'documents/'
},
'credentials': {
'access_key': 'YOUR_ACCESS_KEY',
'secret_key': 'YOUR_SECRET_KEY' # pragma: allowlist secret
},
'job_id': 'my-job-123',
'job_run_id': 'run-456',
'max_files': 100, # Optional: limit number of files
'include_filter': 'pdf,txt,docx', # Optional: file extensions to include
'exclude_filter': 'tmp,log', # Optional: file extensions to exclude
'force_ingest': False # Optional: re-ingest previously processed docs
}
# Create operator instance
ingest_node = IngestSourceOperator(node_config)
# Execute ingestion (input_table is used as trigger)
input_table = pa.Table.from_arrays([])
output_tables, metadata = ingest_node.transform(input_table)
# Access results
result_table = output_tables[0]
print(f"Status: {metadata['node_status']}")
print(f"Documents processed: {metadata['processed_docs']}")
print(f"Total documents: {metadata['total_docs_count']}")
print(f"Failed: {metadata['failed_docs_count']}")
print(f"Skipped: {metadata['skipped_docs_count']}")
print(f"Schema: {result_table.schema}")
See Output Columns below for the full column reference.
See the Usage section above and per-provider configuration in Supported Providers.
| Parameter | Type | Required | Default | Description |
|---|---|---|---|---|
provider |
string | Yes | — | Source provider: filesystem, s3, ibm_cos, sharepoint, onedrive, google_drive, box_driver, dropbox, web |
connection_params |
object | Yes | — | Provider-specific connection settings |
credentials |
object | Yes | — | Provider-specific authentication credentials |
include_filter |
string | No | all types | Comma-separated file extensions to include (no dot) |
max_files |
integer | No | unlimited | Maximum files to ingest |
force_ingest |
boolean | No | false |
Re-ingest previously processed files |
retain_deleted_docs |
boolean | No | false |
Keep records of deleted files |
This operator produces a new table; it does not receive an input table.
| Column | PyArrow Type | Description |
|---|---|---|
id |
string |
Document ID (MD5 hash of source path) |
name |
string |
Source path/identifier |
document_format |
string |
File extension (e.g. .pdf, .docx) |
metadata |
string |
JSON-serialised metadata from the source |
source_id |
string |
Source identifier (from metadata.source) |
path |
string |
Source path/URL used for on-demand binary loading by ExtractOperator |
modified_time |
int64 |
Document modification timestamp (Unix epoch) |
{
"type": "ingest_source",
"name": "ingest_s3_documents",
"config": {
"provider": "s3",
"connection_params": {
"bucket": "my-bucket",
"prefix": "documents/"
},
"credentials": {
"aws_access_key_id": "${AWS_ACCESS_KEY_ID}",
"aws_secret_access_key": "${AWS_SECRET_ACCESS_KEY}"
},
"include_filter": "pdf,docx",
"max_files": 500
}
}
# Convert to pandas for analysis
df = result_table.to_pandas()
# Access individual documents
for i in range(result_table.num_rows):
text = result_table['text'][i].as_py()
metadata = json.loads(result_table['metadata'][i].as_py())
source = result_table['source_id'][i].as_py()
print(f"Document {i+1}:")
print(f" Source: {source}")
print(f" Text length: {len(text)}")
print(f" Metadata: {metadata}")
User Configuration (Flow JSON):
include_filter and exclude_filter parameters (comma-separated strings)"include_filter": "pdf,docx,txt"Internal Implementation (For Connector Developers):
included_extensions (list) → file_extensions (config model field)included_extensions or file_extensions in their flow configurationsThe operator validates and filters files by extension using centralized constants from OperatorConstants.FileExtensions:
Supported Extensions:
Filter Parameters:
Validation Behavior:
include_filter or exclude_filter raise ValueErrorThe operator automatically filters out:
/. (except . and ..)Use the max_files parameter to limit the number of documents processed (default: 100).
This ensures only actual file content is processed, improving efficiency and data quality.
The operator supports incremental processing to avoid re-ingesting unchanged documents:
node_config = {
'provider': 's3',
'connection_params': {...},
'credentials': {...},
'job_id': 'my-job-123',
'force_ingest': False # Set to True to re-ingest all documents
}
Documents are tracked by their ID and modification time. Previously processed documents are automatically skipped unless force_ingest is set to True.
The operator handles errors gracefully following the AbstractOperator pattern:
# Metadata structure (follows AbstractOperator pattern):
metadata = {
"node_status": "completed" | "completed_with_errors" | "completed_with_warnings",
"total_docs_count": 100,
"processed_docs": 95,
"failed_docs_count": 3,
"failed_docs": [
{"id": "doc1", "name": "file1.pdf", "reason": "Error description", "document_url": ""}
],
"skipped_docs_count": 2,
"skipped_docs": [
{"id": "doc2", "name": "file2.pdf", "reason": "Already processed", "document_url": ""}
]
}
Authentication Errors:
Error: Invalid credentials
Solution: Verify credentials are correct and have necessary permissions.
Google Drive Scope Errors:
Error: ('invalid_scope: Bad Request', {'error': 'invalid_scope', 'error_description': 'Bad Request'})
Solution: This error occurs when OAuth scopes are missing or incorrect. To fix:
scopes parameter is included in credentials configuration~/.credentials/token.json)['https://www.googleapis.com/auth/drive.readonly'] for read-only accessConnection Errors:
Error: Could not connect to endpoint
Solution: Check network connectivity and endpoint URLs (especially for S3-compatible storage like IBM COS).
Permission Errors:
Error: Access denied
Solution: Ensure credentials have read permissions for the specified resources.
recursive=False for large folder structuresprefix parameter to limit scopeThe output format is designed for seamless integration with:
# 1. Ingest documents
ingest_node = IngestSourceOperator(ingest_config)
tables, metadata = ingest_node.transform(input_table)
# 2. Process with downstream operators
# embedding_node = EmbeddingOperator(embedding_config)
# embedded_tables, _ = embedding_node.transform(tables[0])
# 3. Store results
# storage_node = StorageOperator(storage_config)
# storage_node.transform(embedded_tables[0])
Enable detailed logging by examining the operator output:
output_tables, metadata = ingest_node.transform(input_table)
print(f"Metadata: {metadata}")
if metadata['status'] == 'error':
print(f"Error: {metadata['message']}")
Test each provider independently:
# Test S3 connectivity
import boto3
s3_client = boto3.client('s3',
aws_access_key_id='YOUR_KEY',
aws_secret_access_key='YOUR_SECRET') # pragma: allowlist secret
response = s3_client.list_objects_v2(Bucket='your-bucket', MaxKeys=1)
print(f"Connection successful: {response['ResponseMetadata']['HTTPStatusCode'] == 200}")
Issue: No documents loaded
Issue: Metadata parsing errors
Issue: Slow performance
langchain==1.2.10
langchain-core==1.2.14
pyarrow==24.0.0
pandas==2.3.3
botocore==1.42.55
boto3==1.42.55, langchain-community==0.4.1google-auth-oauthlib==1.2.4, google-auth-httplib2==0.3.0, google-api-python-client==2.190.0, langchain-google-community==3.0.5O365==2.1.9, langchain-community==0.4.1box-sdk-gen==1.17.0, langchain-community==0.4.1dropbox==12.2.1pypdf2==3.0.1, unstructured[pdf]>=0.10.0google-cloud-storage==3.9.0azure-storage-blob==12.28.0Using uv (recommended):
# From project root directory
# Core installation (includes langchain and langchain-core)
uv sync
# AWS/S3 support
uv sync --extra aws
# Google Drive support
uv sync --extra google-drive
# Microsoft (SharePoint/OneDrive) support
uv sync --extra microsoft
# Box support
uv sync --extra box
# All cloud providers
uv sync --extra all-cloud
# Development dependencies
uv sync --extra dev
Using pip:
# Core installation
pip install langchain==1.2.10 langchain-core==1.2.14 pyarrow==24.0.0 pandas==2.3.3
# AWS/S3 support
pip install boto3==1.42.55 langchain-community==0.4.1
# Google Drive support
pip install google-auth-oauthlib==1.2.4 google-auth-httplib2==0.3.0 google-api-python-client==2.190.0 langchain-google-community==3.0.5 pypdf2==3.0.1 "unstructured[pdf]>=0.10.0"
# Microsoft support
pip install O365==2.1.9 langchain-community==0.4.1
# Box support
pip install box-sdk-gen==1.17.0 langchain-community==0.4.1
Inherits from: AbstractOperator
__init__(node_config: dict)Initialize the operator with configuration.
Parameters:
node_config (dict): Configuration dictionary containing:
provider (str): Provider identifier (s3, google_drive, sharepoint, onedrive, box_driver, dropbox, filesystem, web, custom)connection_params (dict): Provider-specific connection parameterscredentials (dict): Authentication credentialsjob_id (str, optional): Job identifier for trackingjob_run_id (str, optional): Job run identifiermax_files (int, optional): Maximum number of files to process (default: 100)include_filter (str, optional): Comma-separated file extensions to include (defaults to all supported extensions if not specified; must be subset of supported extensions)exclude_filter (str, optional): Comma-separated file extensions to exclude (must be subset of supported extensions)force_ingest (bool, optional): Force re-ingestion of previously processed documents (default: False)Raises:
ValueError: If include_filter or exclude_filter contain unsupported file extensionstransform(input_table: pa.Table) -> tuple[list[pa.Table], dict]Execute document ingestion.
Parameters:
input_table (pa.Table): Input PyArrow table (can be None for initial ingestion)Returns:
tuple[list[pa.Table], dict]:
Raises:
get_metadata() -> dictGet operator metadata including features and attributes.
Returns:
dict: Operator metadata with features, attributes, and availability informationnode_config = {
'provider': 'filesystem',
'connection_params': {
'paths': ['/data/customer_support_docs'],
'recursive': True,
'exclude_patterns': ['*.tmp', '__pycache__/*'],
'max_file_size_mb': 100,
'follow_symlinks': False
},
'include_filter': 'pdf,docx,txt',
'max_files': 500
}
node_config = {
'provider': 'filesystem',
'connection_params': {
'paths': [
'/data/invoices',
'/data/contracts',
'/data/reports'
],
'recursive': True,
'exclude_patterns': ['*.tmp'],
'max_file_size_mb': 50,
'follow_symlinks': False
},
'include_filter': 'pdf,docx',
'force_ingest': False
}
{
"name": "ingest",
"type": "ingest_source",
"config": {
"provider": "filesystem",
"connection_params": {
"paths": [
"./data/invoices",
"./data/contracts"
],
"recursive": true,
"exclude_patterns": ["*.tmp", "__pycache__/*"],
"max_file_size_mb": 100,
"follow_symlinks": false
},
"include_filter": "pdf,docx,txt",
"max_files": 1000,
"force_ingest": false
}
}
node_config = {
'provider': 's3',
'connection_params': {
'bucket': 'company-documents',
'prefix': '2024/invoices/' # Ingests all files in this folder
},
'credentials': {
'access_key': os.getenv('AWS_ACCESS_KEY'),
'secret_key': os.getenv('AWS_SECRET_KEY')
}
}
node_config = {
'provider': 's3',
'connection_params': {
'bucket': 'company-documents',
'prefix': '2024/invoices/report.pdf' # Ingests only this specific file
},
'credentials': {
'access_key': os.getenv('AWS_ACCESS_KEY'),
'secret_key': os.getenv('AWS_SECRET_KEY')
}
}
node_config = {
'provider': 's3',
'connection_params': {
'bucket': 'enterprise-data',
'prefix': 'contracts/',
'endpoint_url': 'https://s3.eu-gb.cloud-object-storage.appdomain.cloud'
},
'credentials': {
'access_key': os.getenv('IBM_COS_ACCESS_KEY'),
'secret_key': os.getenv('IBM_COS_SECRET_KEY')
}
}
node_config = {
'provider': 'google_drive',
'connection_params': {
'folder_id': '1DKN_mxnoW1Uaacghz8vyEeqw-j4IOSFK',
'recursive': True
},
'credentials': {
'credentials_json_path': os.getenv('GOOGLE_CREDENTIALS_PATH'),
'token_path': os.path.expanduser('~/.credentials/gdrive_token.json'),
'scopes': ['https://www.googleapis.com/auth/drive.readonly']
}
}
node_config = {
'provider': 'box_driver',
'connection_params': {
'folder_id': '123456789', # Specific Box folder ID (use '0' for root)
'recursive': True,
'max_file_size_mb': 50,
'exclude_patterns': ['*.tmp', 'Trash/*']
},
'credentials': {
'credentials_json_path': os.getenv('BOX_JWT_CONFIG_FILE')
},
'include_filter': 'pdf,docx,txt,pptx,xlsx', # File extensions to include
'max_files': 100
}
node_config = {
'provider': 'dropbox',
'connection_params': {
'folder_path': '/Reports/2026',
'recursive': True,
'max_file_size_mb': 50,
'exclude_patterns': ['*.tmp', '*/Archive/*']
},
'credentials': {
'access_token': os.getenv('DROPBOX_ACCESS_TOKEN')
},
'include_filter': 'pdf,docx,txt',
'max_files': 100
}
To add support for a new provider:
_get_loader() methodSee sample_flows/use_cases/s3_to_opensearch.json for a complete example ingesting from S3 through to OpenSearch.
See project LICENSE file for details.