Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 7 additions & 0 deletions backend/connector_v2/unstract_account.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,9 +2,12 @@
import os

import boto3
from botocore.config import Config
from botocore.exceptions import ClientError
from django.conf import settings

from unstract.sdk1.patches.storage_compat import S3_CHECKSUM_CONFIG

logger = logging.getLogger(__name__)


Expand All @@ -24,6 +27,8 @@ def provision_s3_storage(self) -> None:
aws_access_key_id=access_key,
aws_secret_access_key=secret_key,
endpoint_url="https://storage.googleapis.com",
# GCS S3 interop: no aws-chunked uploads (UN-4224).
config=Config(**S3_CHECKSUM_CONFIG),
)

# Check if folder exists and create if it is not available
Expand Down Expand Up @@ -53,6 +58,8 @@ def upload_sample_files(self) -> None:
aws_access_key_id=access_key,
aws_secret_access_key=secret_key,
endpoint_url="https://storage.googleapis.com",
# GCS S3 interop: no aws-chunked uploads (UN-4224).
config=Config(**S3_CHECKSUM_CONFIG),
)

folder = f"{self.tenant}/{self.username}/input/examples/"
Expand Down
2 changes: 1 addition & 1 deletion backend/pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -14,7 +14,7 @@ classifiers = [

dependencies = [
"Authlib==1.6.12", # For Auth plugins
"boto3~=1.34.0", # For Unstract-cloud-storage
"boto3==1.43.106", # For Unstract-cloud-storage
Comment thread
praveen-formido marked this conversation as resolved.
"celery[amqp]>=5.3.4", # For Celery
"cron-descriptor==1.4.0", # For cron string description
"cryptography>=48.0.1",
Expand Down
2,139 changes: 1,076 additions & 1,063 deletions backend/uv.lock

Large diffs are not rendered by default.

1,701 changes: 876 additions & 825 deletions platform-service/uv.lock

Large diffs are not rendered by default.

10 changes: 5 additions & 5 deletions unstract/connectors/pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -25,15 +25,15 @@ dependencies = [
# google-auth >= 2.53 no longer depends on urllib3 (2.20 did) — that's what unblocks urllib3 2.x
"google-auth>=2.22.0,<3",
"google-cloud-secret-manager==2.16.1",
"google-cloud-storage==2.9.0",
"google-cloud-storage==3.16.0",
# Filesystem connectors
"s3fs[boto3]==2024.10.0", # For Minio
"s3fs==2026.9.0", # For Minio
"PyDrive2[fsspec]==1.15.4", # For GDrive
"oauth2client==4.1.3", # For GDrive
"dropboxdrivefs==1.4.1", # For Dropbox
"boxfs==0.2.1", # For Box
"gcsfs==2024.10.0", # For GoogleCloudStorage
"adlfs~=2024.7.0", # For AzureCloudStorage
"gcsfs==2026.10.0", # For GoogleCloudStorage
Comment thread
praveen-formido marked this conversation as resolved.
"adlfs~=2026.8.0", # For AzureCloudStorage
"Office365-REST-Python-Client~=2.6.0", # For SharePoint/OneDrive
# Database connectors
"psycopg2-binary==2.9.9", # For Postgres, Redshift
Expand All @@ -42,7 +42,7 @@ dependencies = [
"pymssql==2.3.4", # For MSSQL
"PyMySQL==1.1.1", # For MySQL
"oracledb==2.4.0", # For OracleDB
"fsspec[sftp]~=2024.10.0", # For SFTP
"fsspec[sftp]~=2026.9.0", # For SFTP
]

# [build-system]
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,10 @@

from fsspec import AbstractFileSystem

# Keeps gcsfs on its core GCSFileSystem (gcsfs 2026.x otherwise swaps in the
# experimental gRPC-backed ExtendedGcsFileSystem); must run before gcsfs is
# imported below (UN-4224).
import unstract.sdk1.patches.storage_compat # noqa: F401
from unstract.connectors.exceptions import ConnectorError
from unstract.connectors.filesystems.unstract_file_system import UnstractFileSystem

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,10 @@
from unstract.connectors.exceptions import ConnectorError
from unstract.connectors.filesystems.unstract_file_system import UnstractFileSystem

# sdk1 is already a runtime dependency here via `unstract.filesystem`. Importing
# storage_compat also re-adds Content-MD5 on DeleteObjects (UN-4224).
from unstract.sdk1.patches.storage_compat import S3_CHECKSUM_CONFIG

from .exceptions import (
BUCKET_PROBE_DISPOSITION,
BucketProbeDisposition,
Expand All @@ -24,36 +28,6 @@
logger = logging.getLogger(__name__)


class _BucketScopedFileSystem(DirFileSystem):
"""`DirFileSystem.walk()` relpaths the directory string it yields, but
not the `name` field inside each file/dir entry's own metadata dict —
those still carry the wrapped fs's raw, bucket-prefixed key. `ls()`
already fixes every entry; `walk()` doesn't. Fix it here so discovery
(which walks) and browsing (which lists) agree (UN-3487).
"""

def _relpath_entries(
self, entries: dict[str, Any] | list[str]
) -> dict[str, Any] | list[str]:
# detail=False (fsspec's own default) yields bare basenames with
# nothing to fix. detail=True yields a dict already keyed by bare
# basename — only each entry's own `name` field is bucket-qualified.
if not isinstance(entries, dict):
return entries
return {
name: {**info, "name": self._relpath(info["name"])}
for name, info in entries.items()
}

def walk(self, path: str, *args: Any, **kwargs: Any) -> Any:
for root, dirs, files in super().walk(path, *args, **kwargs):
yield root, self._relpath_entries(dirs), self._relpath_entries(files)

async def _walk(self, path: str, *args: Any, **kwargs: Any) -> Any:
async for root, dirs, files in super()._walk(path, *args, **kwargs):
yield root, self._relpath_entries(dirs), self._relpath_entries(files)


# Cap concurrent per-bucket probes to avoid S3 503 SlowDown on large accounts.
_MAX_CONCURRENT_BUCKET_PROBES = 16
_BUCKET_PROBE_RETRY_DELAY_SECONDS = 0.5
Expand Down Expand Up @@ -205,6 +179,9 @@ def __init__(self, settings: dict[str, Any]):
default_cache_type="none",
skip_instance_cache=True,
client_kwargs=client_kwargs,
# No aws-chunked uploads: UCS (GCS's S3 API) and TLS MinIO get the
# plain payload-signed requests boto3 1.34 sent (UN-4224).
config_kwargs=dict(S3_CHECKSUM_CONFIG),
**creds,
)

Expand Down Expand Up @@ -371,12 +348,13 @@ def get_fsspec_fs(self) -> AbstractFileSystem:
"""Return the filesystem scoped to this connector's bucket.

When a bucket is configured, every operation (list, read, write,
`test_credentials`) is confined to it via `_BucketScopedFileSystem` —
the underlying credentials may see more, but this connector never
will.
`test_credentials`) is confined to it via `DirFileSystem` — the
underlying credentials may see more, but this connector never will.
Since fsspec 2026.x, `DirFileSystem.walk()` also relpaths each entry's
own `name` field, so walk and ls agree without an override (UN-3487).
"""
if self.bucket:
return _BucketScopedFileSystem(path=self.bucket, fs=self.s3)
return DirFileSystem(path=self.bucket, fs=self.s3)
return self.s3

def test_credentials(self) -> bool:
Expand Down
11 changes: 6 additions & 5 deletions unstract/connectors/tests/filesystems/test_miniofs.py
Original file line number Diff line number Diff line change
Expand Up @@ -278,12 +278,13 @@ def test_get_fsspec_fs_scopes_to_bucket(self) -> None:
self.assertIs(scoped.fs, fs.s3)

def test_walk_results_are_relative_to_the_bucket(self) -> None:
# `DirFileSystem.walk()` relpaths the directory string it yields,
# but not the `name` field inside each entry's own metadata dict —
# those still carry the wrapped fs's raw, bucket-prefixed key. This
# is what workflow-execution file discovery reads (it walks, the
# UI browser lists) — a name still carrying the bucket here is what
# Each walk entry's own `name` field must come back bucket-relative,
# not as the wrapped fs's raw, bucket-prefixed key (UN-3487). This is
# what workflow-execution file discovery reads (it walks, the UI
# browser lists) — a name still carrying the bucket here is what
# doubles the prefix at the point a discovered file gets opened.
# fsspec 2026.x's `DirFileSystem.walk()` does this itself; this test
# guards against a regression either upstream or in `MinioFS`.
#
# Matches fsspec's real shape (verified against AbstractFileSystem.
# walk): the dict is keyed by bare basename already; only each
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,64 @@
"""The connectors load sdk1's storage_compat on their own (UN-4224).

Each probe runs in a fresh interpreter and imports only the connector module,
so nothing in this test process can satisfy it. Deleting the import from the
connector must fail these tests.
"""

import os
import subprocess
import sys

from unstract.connectors.filesystems.minio.minio import MinioFS
from unstract.connectors.filesystems.ucs.ucs import UnstractCloudStorage

_EXPECTED_CHECKSUM_CONFIG = {
"request_checksum_calculation": "when_required",
"response_checksum_validation": "when_required",
}


def _run_probe(probe: str) -> None:
env = {
k: v for k, v in os.environ.items() if k != "GCSFS_EXPERIMENTAL_ZB_HNS_SUPPORT"
}
# The connectors package is on the path via pytest's `pythonpath`, which a
# child interpreter does not inherit.
env["PYTHONPATH"] = os.pathsep.join(p for p in sys.path if p)
result = subprocess.run(
[sys.executable, "-c", probe], env=env, capture_output=True, text=True
)
assert result.returncode == 0, result.stderr[-2000:]


def test_minio_connector_import_registers_delete_objects_md5() -> None:
_run_probe(
"import unstract.connectors.filesystems.minio.minio\n"
"import botocore.handlers\n"
"names = [h[1].__name__ for h in botocore.handlers.BUILTIN_HANDLERS\n"
" if h[0] == 'before-call.s3.DeleteObjects']\n"
"assert 'add_content_md5' in names, names\n"
)


def test_gcs_connector_import_keeps_core_gcsfs() -> None:
_run_probe(
"import unstract.connectors.filesystems.google_cloud_storage."
"google_cloud_storage\n"
"import gcsfs\n"
"assert gcsfs.GCSFileSystem.__name__ == 'GCSFileSystem', gcsfs.GCSFileSystem\n"
)


def test_minio_and_ucs_send_plain_uploads() -> None:
# No aws-chunked uploads: UCS is GCS's S3 interop API, MinIO may sit behind
# TLS; both get the payload-signed requests boto3 1.34 sent.
settings = {
"key": "k",
"secret": "s",
"endpoint_url": "https://storage.example",
"bucket": "b",
}
for fs_class in (MinioFS, UnstractCloudStorage):
s3 = fs_class(settings).s3
assert s3.config_kwargs == _EXPECTED_CHECKSUM_CONFIG, fs_class.__name__
Loading
Loading