Skip to content

Commit 382b164

Browse files
committed
added acl support
1 parent e6db48f commit 382b164

3 files changed

Lines changed: 27 additions & 1 deletion

File tree

src/api/router.py

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -54,6 +54,7 @@ async def upload_ingest_router(
5454
session_manager=session_manager,
5555
task_service=task_service,
5656
user=user,
57+
settings_json=settings_json,
5758
)
5859

5960
logger.debug("Routing to Langflow upload-ingest pipeline via task service")
@@ -78,12 +79,20 @@ async def _traditional_upload_ingest_task(
7879
session_manager,
7980
task_service,
8081
user: User,
82+
settings_json: str | None = None,
8183
):
8284
"""Task-based traditional upload and ingest for single/multiple files"""
8385
try:
8486
if not upload_files:
8587
return JSONResponse({"error": "Missing files"}, status_code=400)
8688

89+
settings = None
90+
if settings_json:
91+
try:
92+
settings = json.loads(settings_json)
93+
except json.JSONDecodeError as e:
94+
return JSONResponse({"error": f"Invalid settings JSON: {e}"}, status_code=400)
95+
8796
user_id = user.user_id
8897
user_name = user.name
8998
user_email = user.email
@@ -121,6 +130,7 @@ async def _traditional_upload_ingest_task(
121130
owner_email=user_email,
122131
original_filenames=file_path_to_original_filename,
123132
replace_duplicates=replace_duplicates,
133+
settings=settings,
124134
)
125135

126136
return JSONResponse(

src/models/processors.py

Lines changed: 15 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -190,7 +190,7 @@ async def process_document_standard(
190190
chunk_size: int = None,
191191
chunk_overlap: int = None,
192192
is_sample_data: bool = False,
193-
acl: "DocumentACL" = None,
193+
acl: "DocumentACL | None" = None,
194194
):
195195
"""
196196
Standard processing pipeline for non-Langflow processors:
@@ -404,6 +404,7 @@ def __init__(
404404
docling_service=None,
405405
replace_duplicates: bool = False,
406406
session_manager=None,
407+
settings: dict | None = None,
407408
):
408409
super().__init__(
409410
document_service,
@@ -421,6 +422,7 @@ def __init__(
421422
self.session_manager = session_manager or (
422423
document_service.session_manager if document_service else None
423424
)
425+
self.settings = settings
424426
if self.session_manager is None:
425427
raise ValueError("session_manager is required for DocumentFileProcessor")
426428

@@ -480,6 +482,17 @@ async def process_item(self, upload_task: UploadTask, item: str, file_task: File
480482
except Exception:
481483
file_size = 0
482484

485+
# Parse ACL from settings if present
486+
from connectors.base import DocumentACL
487+
488+
acl = None
489+
if self.settings and (self.settings.get("allowed_users") is not None or self.settings.get("allowed_groups") is not None):
490+
acl = DocumentACL(
491+
owner=self.owner_user_id,
492+
allowed_users=self.settings.get("allowed_users", []),
493+
allowed_groups=self.settings.get("allowed_groups", [])
494+
)
495+
483496
# Use consolidated standard processing
484497
result = await self.process_document_standard(
485498
file_path=item,
@@ -492,6 +505,7 @@ async def process_item(self, upload_task: UploadTask, item: str, file_task: File
492505
file_size=file_size,
493506
connector_type=self.connector_type,
494507
is_sample_data=self.is_sample_data,
508+
acl=acl,
495509
)
496510

497511
file_task.status = TaskStatus.COMPLETED

src/services/task_service.py

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -129,6 +129,7 @@ async def create_upload_task(
129129
owner_email: str = None,
130130
original_filenames: dict | None = None,
131131
replace_duplicates: bool = False,
132+
settings: dict | None = None,
132133
) -> str:
133134
"""Create a new upload task for bulk file processing"""
134135
# Use default DocumentFileProcessor with user context
@@ -144,6 +145,7 @@ async def create_upload_task(
144145
docling_service=self.docling_service,
145146
replace_duplicates=replace_duplicates,
146147
session_manager=self.session_manager,
148+
settings=settings,
147149
)
148150
return await self.create_custom_task(
149151
user_id,

0 commit comments

Comments
 (0)