docling-pipelines

IngestSourceOperator

Ingests document metadata from cloud storage and collaboration platforms (S3, SharePoint, OneDrive, Google Drive, Box, Dropbox).


Overview

IngestSourceOperator 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.


Key Features

Supported Providers

1. Local Filesystem

Ingest 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").

2. Amazon S3 and S3-Compatible Storage

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:

Note: File extension filtering is configured at the operator level using include_filter and exclude_filter parameters (see File Filtering section below).

3. Microsoft SharePoint

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:

Parameters:

4. Microsoft OneDrive

Ingest 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:

Parameters:

5. Google Drive

Ingest 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:

Parameters:

OAuth Scopes: The operator uses read-only access by default for security. Available scopes:

Important: If you change scopes, you must delete the existing token file to re-authenticate with the new permissions.

6. Box

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:

Parameters:

Box JWT Setup:

  1. Create a Box application in the Box Developer Console
  2. Choose “Server Authentication (with JWT)” as authentication method
  3. Configure application permissions:
    • Read all files and folders stored in Box
    • Manage enterprise properties
  4. Generate a public/private keypair
  5. Download the JWT configuration JSON file
  6. Submit application for admin approval (if required)
  7. Admin must authorize the application in Box Admin Console

Authentication Flow: The adapter uses JWT (JSON Web Token) authentication which provides:

Security Notes:

7. Dropbox

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:

Parameters:

Dropbox App Setup:

  1. Create a scoped app in the Dropbox App Console
  2. Choose App folder access to limit the app to a single folder, or Full Dropbox
  3. On the Permissions tab enable account_info.read, files.metadata.read and files.content.read, then submit
  4. Generate an access token on the Settings tab, or run the OAuth2 flow with token_access_type=offline to obtain a refresh token
  5. Export the token as an environment variable and reference it as ${DROPBOX_ACCESS_TOKEN} in the flow

Authentication Notes:

Security Notes:

8. Custom Loaders

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:

Requirements:

Usage

Basic Example

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}")

Output Schema

See Output Columns below for the full column reference.


Parameters

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

Output Columns

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)

Operator Configuration

{
  "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
  }
}

Accessing Results

# 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}")

Parameter Naming Clarification

User Configuration (Flow JSON):

Internal Implementation (For Connector Developers):

File Filtering

Extension-Based Filtering

The operator validates and filters files by extension using centralized constants from OperatorConstants.FileExtensions:

Supported Extensions:

Filter Parameters:

Validation Behavior:

S3 Filtering

The operator automatically filters out:

Max Files Limit

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.

Incremental Updates

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.

Error Handling

Graceful Degradation

The operator handles errors gracefully following the AbstractOperator pattern:

Metadata Response

# 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": ""}
    ]
}

Common Errors

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:

  1. Ensure the scopes parameter is included in credentials configuration
  2. Delete the existing token file (default: ~/.credentials/token.json)
  3. Re-run the ingestion to trigger re-authentication with correct scopes
  4. Use the default scope ['https://www.googleapis.com/auth/drive.readonly'] for read-only access

Connection 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.

Performance Considerations

Large Datasets

Memory Usage

Optimization Tips

  1. Use specific prefixes/folder IDs to limit scope
  2. Filter file types at the source when possible
  3. Process in batches for very large datasets
  4. Use appropriate loader configurations for your use case

Integration with Downstream Operators

The output format is designed for seamless integration with:

Example Pipeline

# 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])

Security Best Practices

  1. Credential Management:
    • Never hardcode credentials in source code
    • Use environment variables or secret management systems
    • Rotate credentials regularly
  2. Access Control:
    • Use least-privilege principle for service accounts
    • Limit bucket/folder access to necessary resources
    • Monitor access logs for suspicious activity
  3. Data Protection:
    • Use encrypted connections (HTTPS/TLS)
    • Consider encrypting sensitive data at rest
    • Implement data retention policies
  4. OAuth Tokens (Google Drive, SharePoint, OneDrive):
    • Store tokens securely with restricted file permissions
    • Never commit token files to version control
    • Implement token refresh mechanisms

Troubleshooting

Debug Mode

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']}")

Testing Connectivity

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}")

Common Issues

Issue: No documents loaded

Issue: Metadata parsing errors

Issue: Slow performance

Dependencies

Core Dependencies

langchain==1.2.10
langchain-core==1.2.14
pyarrow==24.0.0
pandas==2.3.3
botocore==1.42.55

Provider-Specific Dependencies

Installation

Using 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

API Reference

Class: IngestSourceOperator

Inherits from: AbstractOperator

__init__(node_config: dict)

Initialize the operator with configuration.

Parameters:

Raises:

transform(input_table: pa.Table) -> tuple[list[pa.Table], dict]

Execute document ingestion.

Parameters:

Returns:

Raises:

get_metadata() -> dict

Get operator metadata including features and attributes.

Returns:

Examples

Example 1: Filesystem — Single Directory (Python)

node_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
}

Example 2: Filesystem — Multiple Directories (Python)

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
}

Example 3: Filesystem — Flow JSON (Multiple Directories)

{
  "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
  }
}

Example 4: S3 with Folder Prefix Filtering

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')
    }
}

Example 5: S3 with File-Level Ingestion

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')
    }
}

Example 6: S3-Compatible Storage (IBM COS)

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')
    }
}

Example 7: Google Drive Recursive

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']
    }
}

Example 8: Box with JWT Authentication

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
}

Example 9: Dropbox Folder with Extension Filtering

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
}

Contributing

To add support for a new provider:

  1. Add the provider to the _get_loader() method
  2. Import the corresponding LangChain loader
  3. Map configuration parameters to loader initialization
  4. Update this documentation with provider details
  5. Add example configuration and usage

Sample Flow

See sample_flows/use_cases/s3_to_opensearch.json for a complete example ingesting from S3 through to OpenSearch.

License

See project LICENSE file for details.