import base64 import os from datetime import timedelta, datetime from io import BytesIO import logging from django.core.exceptions import ObjectDoesNotExist from minio.error import S3Error import chardet import fitz import unicodedata from cfgv import ValidationError from django.core.files.base import ContentFile from django.db import transaction from minio import Minio, S3Error from minio.commonconfig import CopySource from DCbackend import settings from DCbackend.utils.common import success, fail from backend.models import Admin, Knowledgebase, File, DocumentKbm, File2document from pypinyin import lazy_pinyin, Style import re # 配置 logger logger = logging.getLogger(__name__) # 初始化 MinIO 客户端 minio_client = Minio( settings.MINIO_ENDPOINT, access_key=settings.MINIO_ACCESS_KEY, secret_key=settings.MINIO_SECRET_KEY, secure=settings.MINIO_SECURE ) def contains_chinese(text): """检查字符串是否包含中文字符""" return any('\u4e00' <= char <= '\u9fff' for char in text) def convert_to_pinyin(text): """ 将中文字符转换为不带声调的拼音,保留字母和数字,其他字符替换为连字符, 并确保生成的名称符合 bucket 命名规则 """ result = [] for char in text: if '\u4e00' <= char <= '\u9fff': # 中文字符 pinyin = lazy_pinyin(char, style=Style.NORMAL) result.extend(pinyin) elif char.isalnum(): # 字母和数字 result.append(char.lower()) else: # 其他字符替换为连字符 result.append('-') # 合并字符并去除声调 bucket_name = ''.join(result) bucket_name = unicodedata.normalize('NFKD', bucket_name).encode('ASCII', 'ignore').decode('ASCII') # 合并连续的连字符 bucket_name = re.sub(r'-+', '-', bucket_name) # 确保名称以字母或数字开头和结尾 bucket_name = bucket_name.strip('-') # 如果名称为空,使用默认名称 if not bucket_name: bucket_name = 'default-bucket' # 确保名称长度在 3-63 之间 if len(bucket_name) < 3: bucket_name = bucket_name.ljust(3, 'a') elif len(bucket_name) > 63: bucket_name = bucket_name[:63] return bucket_name.lower() class MinioService: @staticmethod def listBuckets(request): #查询所有bucket信息 buckets = minio_client.list_buckets() # 将 Bucket 对象转换为可序列化的字典,包含所有可用信息 serializable_buckets = [] for bucket in buckets: # 处理创建日期 if isinstance(bucket.creation_date, datetime): # 添加8小时 adjusted_date = bucket.creation_date + timedelta(hours=8) formatted_date = adjusted_date.strftime("%Y-%m-%d %H:%M:%S") else: formatted_date = str(bucket.creation_date) bucket_info = { 'bucket_name': bucket.name, 'creation_date': formatted_date, } #根据桶名查询数量 try: objects = minio_client.list_objects(bucket.name, recursive=True) file_count = sum(1 for _ in objects) bucket_info['file_count'] = file_count except Exception as e: bucket_info['file_count'] = f"Error: {str(e)}" serializable_buckets.append(bucket_info) return success(serializable_buckets) @staticmethod def createBucket(request): """创建新的 bucket""" try: user_id = request.POST.get("user_id") # 从请求中获取 bucket 名称 bucket_name = request.POST.get('bucket_name') description = request.POST.get('description', "") if not bucket_name: return fail('Bucket name is required') minio_name = bucket_name if contains_chinese(bucket_name): minio_name = convert_to_pinyin(bucket_name) # 检查 bucket 是否已存在 if minio_client.bucket_exists(minio_name): return fail(f'Bucket "{minio_name}" 已存在') # 创建 bucket minio_client.make_bucket(minio_name) # 同步数据库 return MinioService.saveBucketDb(bucket_name, user_id, minio_name, description) except S3Error as e: return fail(f'Failed to create bucket: {str(e)}') except Exception as e: return fail(f'An unexpected error occurred: {str(e)}') @staticmethod def saveBucketDb(bucket_name, user_id,minio_name,description): try: admin = Admin.objects.get(id=user_id) except Admin.DoesNotExist: return fail("用户信息不存在") try: # 假设 Knowledgebase 是您的模型类名 db = Knowledgebase( role_id=admin.role_id, name=bucket_name, location=minio_name, description=description, created_by=user_id ) db.save() file = File.objects.filter( role_id=admin.role_id, created_by=user_id, name='/', source_type='' ).exclude(status=4).first() if file is None: # 如果记录不存在,创建新记录 file = File.objects.create( role_id=admin.role_id, name='/', source_type='', created_by= user_id, location= '', type= "folder", status= 5, ) # 创建后更新parent_id file.parent_id = file.id file.save() kbm = File.objects.create( role_id=admin.role_id, name='.knowledgebase', # 根据需要修改名称 source_type='knowledgebase', created_by=user_id, location='', type="folder", status=5, parent_id=file.id # 设置parent_id为主记录的id ) bucket = File.objects.create( role_id=admin.role_id, name=bucket_name, # 根据需要修改名称 source_type='knowledgebase', created_by=user_id, location='', type="folder", status=5, parent_id=kbm.id ) else: kbm = File.objects.get( role_id=admin.role_id, created_by=user_id, name='.knowledgebase', type='folder', source_type='knowledgebase' ) bucket = File.objects.create( role_id=admin.role_id, name=bucket_name, # 根据需要修改名称 source_type='knowledgebase', created_by=user_id, location='', type="folder", status=5, parent_id=kbm.id ) except ValidationError as e: return fail(f"验证错误: {e}") except Exception as e: return fail(f"保存失败: {str(e)}") return success(f'Bucket "{bucket_name}" 创建成功') @staticmethod def is_valid_bucket_name(bucket_name): """检查 bucket 名称是否合法""" import re # bucket 名称必须在 3-63 个字符之间,只能包含小写字母、数字和连字符 pattern = r'^[a-z0-9][a-z0-9\-]{1,61}[a-z0-9]$' return re.match(pattern, bucket_name) is not None #获取指定buck内文件信息 @staticmethod def getBucketContents(request): """获取指定 bucket 内的所有文件信息""" try: # 从请求中获取 bucket 名称 bucket_name = request.POST.get('bucket_name') # 使用 get 方法获取 page 和 page_size,如果不存在则使用默认值 page = int(request.POST.get('page', 1)) page_size = int(request.POST.get('page_size', 10)) print(page, page_size) if not bucket_name: return fail('请求参数为空') # 检查 bucket 是否存在 if not minio_client.bucket_exists(bucket_name): return fail(f'Bucket "{bucket_name}" 不存在') # 获取 bucket 内的所有对象 objects = minio_client.list_objects(bucket_name, recursive=True) # 整理文件信息 file_info = list(objects) # 转换为列表以获得准确的长度 # 计算总数和总页数 total_count = len(file_info) total_pages = (total_count + page_size - 1) // page_size # 确保页码在有效范围内 page = max(1, min(page, total_pages)) # 计算切片的起始和结束索引 start_index = (page - 1) * page_size end_index = min(start_index + page_size, total_count) # 获取当前页的数据 paginated_files = [ { 'object_name': obj.object_name, 'size': obj.size, 'last_modified': obj.last_modified, 'version_id': obj.version_id, 'etag': obj.etag } for obj in file_info[start_index:end_index] ] return success({ 'bucket_name': bucket_name, 'files': paginated_files, 'page': page, 'page_size': page_size, 'total_pages': total_pages, 'total_count': total_count }) except S3Error as e: return fail(f'Failed to get bucket contents: {str(e)}') except Exception as e: return fail(f'An unexpected error occurred: {str(e)}') @staticmethod @transaction.atomic() # def post(request): # """上传文件""" # uploaded_file = request.FILES['file'] # 获取上传的文件 # bucket_id = request.POST.get('bucket_id') # BUCKET的名称 # user_id = request.POST.get('user_id') # BUCKET的名称 # file_path = request.POST.get('file_path', '') # doc_type_id = request.POST.get('doc_type_id',0) # # if not uploaded_file: # return fail('没有需要上传的文件') # if not bucket_id: # return fail('bucket_id为空') # # # 使用 file_path 和文件名构造对象名称 # object_name = f"{file_path.strip('/')}/{uploaded_file.name}".lstrip('/') # count = DocumentKbm.objects.filter(name=object_name, kb_id=bucket_id).exclude(status=4).count() # if count > 0: # return fail("已有重复文件,请删除后重试") # # try: # # 读取文件内容 # file_content = uploaded_file.read() # # # 使用 MinIO 客户端上传文件 # knowledgebase = Knowledgebase.objects.get(id=bucket_id) # bucket_name = knowledgebase.location # # minio_client.put_object( # bucket_name, # object_name, # ContentFile(file_content), # length=len(file_content), # content_type=uploaded_file.content_type # ) # MinioService.saveDocumentKbm(knowledgebase, uploaded_file, user_id, object_name,doc_type_id) # # return success("保存成功") # except Exception as e: # return fail(str(e)) # def post(request): # """上传文件""" # logger.info("Starting file upload process") # try: # uploaded_file = request.FILES['file'] # bucket_id = request.POST.get('bucket_id') # user_id = request.POST.get('user_id') # file_path = request.POST.get('file_path', '') # doc_type_id = request.POST.get('doc_type_id', 0) # # if not uploaded_file: # return fail('没有需要上传的文件') # if not bucket_id: # return fail('bucket_id为空') # # object_name = f"{file_path.strip('/')}/{uploaded_file.name}".lstrip('/') # logger.debug(f"Constructed object name: {object_name}") # # count = DocumentKbm.objects.filter(name=object_name, kb_id=bucket_id).exclude(status=4).count() # if count > 0: # return fail("已有重复文件,请删除后重试") # # file_content = uploaded_file.read() # knowledgebase = Knowledgebase.objects.get(id=bucket_id) # bucket_name = knowledgebase.location # # logger.info(f"Uploading file to MinIO: {object_name} in bucket {bucket_name}") # minio_client.put_object( # bucket_name, # object_name, # ContentFile(file_content), # length=len(file_content), # content_type=uploaded_file.content_type # ) # # logger.info("MinIO upload successful, calling saveDocumentKbm") # MinioService.saveDocumentKbm(knowledgebase, uploaded_file, user_id, object_name, doc_type_id) # # logger.info("File upload process completed successfully") # return success("保存成功") # except Exception as e: # logger.error(f"Error in file upload process: {str(e)}", exc_info=True) # return fail(str(e)) def post(request): """上传文件""" logger.info("Starting file upload process") try: uploaded_file = request.FILES['file'] bucket_id = request.POST.get('bucket_id') user_id = request.POST.get('user_id') file_path = request.POST.get('file_path', '') doc_type_id = request.POST.get('doc_type_id', 0) if not uploaded_file: return fail('没有需要上传的文件') if not bucket_id: return fail('bucket_id为空') object_name = f"{file_path.strip('/')}/{uploaded_file.name}".lstrip('/') logger.debug(f"Constructed object name: {object_name}") count = DocumentKbm.objects.filter(name=object_name, kb_id=bucket_id).exclude(status=4).count() if count > 0: return fail("已有重复文件,请删除后重试") file_content = uploaded_file.read() try: knowledgebase = Knowledgebase.objects.get(id=bucket_id) except ObjectDoesNotExist: return fail(f"知识库 ID {bucket_id} 不存在") bucket_name = knowledgebase.location logger.info(f"Uploading file to MinIO: {object_name} in bucket {bucket_name}") try: minio_client.put_object( bucket_name, object_name, ContentFile(file_content), length=len(file_content), content_type=uploaded_file.content_type ) except S3Error as s3_err: logger.error(f"MinIO upload failed: {str(s3_err)}") return fail(f"文件上传到 MinIO 失败: {str(s3_err)}") logger.info("MinIO upload successful, calling saveDocumentKbm") try: MinioService.saveDocumentKbm(knowledgebase, uploaded_file, user_id, object_name, doc_type_id) except Exception as save_err: logger.error(f"Error in saveDocumentKbm: {str(save_err)}", exc_info=True) # 考虑在这里删除已上传到 MinIO 的文件 return fail(f"保存文档信息失败: {str(save_err)}") logger.info("File upload process completed successfully") return success("保存成功") except KeyError as key_err: logger.error(f"Missing required field: {str(key_err)}") return fail(f"缺少必要的字段: {str(key_err)}") except Exception as e: logger.error(f"Error in file upload process: {str(e)}", exc_info=True) return fail(str(e)) def is_image_file(extension): image_extensions = [ 'jpg', 'jpeg', 'png', 'gif', 'bmp', 'tiff', 'webp', 'svg', 'raw', 'heif', 'heic', 'indd', 'ai', 'eps', 'psd', 'xcf', 'cr2', 'nef', 'orf', 'sr2', 'jfif', 'exif', 'ico', 'tga' ] return extension.lower().strip('.') in image_extensions def is_text_file(extension): text_extensions = [ 'txt', 'pdf', 'doc', 'docx', 'rtf', 'odt', 'xls', 'xlsx', 'csv', 'tsv', 'json', 'xml', 'html', 'htm', 'md', 'markdown', 'tex', 'log', 'ini', 'cfg', 'conf', 'py', 'js', 'css', 'scss', 'less', 'sql', 'php', 'java', 'c', 'cpp', 'h', 'hpp', 'sh', 'bat', 'ps1', 'rb', 'yaml', 'yml', 'toml', 'rst', 'asciidoc', 'ppt', 'pptx', 'odp', 'key', 'pages', 'numbers' ] return extension.lower().strip('.') in text_extensions @staticmethod @transaction.atomic # def saveDocumentKbm(knowledgebase,uploaded_file,user_id,object_name,doc_type_id): # try: # size = uploaded_file.size # _, file_extension = os.path.splitext(uploaded_file.name) # # file_extension 现在包含了文件的后缀名,包括点号(例如 ".txt") # # 如果您不想要点号,可以这样做: # file_extension = file_extension[1:] if file_extension else '' # # 根据文件类型设置 parser_id # if MinioService.is_image_file(file_extension): # parser_id = 'picture' # elif MinioService.is_text_file(file_extension): # parser_id = 'naive' # # documentKbm = DocumentKbm.objects.create( # kb_id=knowledgebase.id, # parser_id=parser_id, # parser_config='{"pages": [[1, 1000000]]}', # type = file_extension, # created_by=user_id, # name=object_name, # location=object_name, # size=size, # doc_type_id=doc_type_id # ) # doc_id = documentKbm.id # # admin = Admin.objects.get(id=user_id) # file = File.objects.filter( # role_id=admin.role_id, # created_by=user_id, # name=knowledgebase.name, # source_type='knowledgebase', # type='folder' # ).exclude(status=4).first() # # fileDB = File.objects.create( # role_id=admin.role_id, # name=object_name, # source_type='knowledgebase', # created_by=user_id, # location=object_name, # type=file_extension, # status=5, # size=size, # parent_id= file.id # ) # file_id = fileDB.id # # file2document = File2document.objects.create( # file_id=file_id, # document_id=doc_id # ) # # # # child_count = File.objects.filter(parent_id=file.id).exclude(status=4).count() # # knowledgebase.doc_num = child_count # knowledgebase.save() # # # except Exception as e: # # return fail(str(e)) @staticmethod @transaction.atomic def saveDocumentKbm(knowledgebase, uploaded_file, user_id, object_name, doc_type_id): logger.info(f"Starting saveDocumentKbm for file: {object_name}") try: size = uploaded_file.size _, file_extension = os.path.splitext(uploaded_file.name) file_extension = file_extension[1:] if file_extension else '' logger.debug(f"File size: {size}, extension: {file_extension}") if MinioService.is_image_file(file_extension): parser_id = 'picture' elif MinioService.is_text_file(file_extension): parser_id = 'naive' else: parser_id = 'unknown' logger.debug(f"Determined parser_id: {parser_id}") documentKbm = DocumentKbm.objects.create( kb_id=knowledgebase.id, parser_id=parser_id, parser_config='{"pages": [[1, 1000000]]}', type=file_extension, created_by=user_id, name=object_name, location=object_name, size=size, doc_type_id=doc_type_id ) logger.info(f"Created DocumentKbm with id: {documentKbm.id}") try: admin = Admin.objects.get(id=user_id) except Admin.DoesNotExist: logger.error(f"Admin with id {user_id} does not exist") raise ValueError(f"Admin with id {user_id} does not exist") logger.debug( f"Searching for File with: role_id={admin.role_id}, created_by={user_id}, name={knowledgebase.name}") file = File.objects.filter( role_id=admin.role_id, created_by=user_id, name=knowledgebase.name, source_type='knowledgebase', type='folder' ).exclude(status=4).first() if not file: logger.warning(f"No matching File found for knowledgebase: {knowledgebase.name}. Creating a new one.") file = File.objects.create( role_id=admin.role_id, name=knowledgebase.name, source_type='knowledgebase', created_by=user_id, type='folder', status=5 ) logger.info(f"Created new File record for knowledgebase: {file.id}") fileDB = File.objects.create( role_id=admin.role_id, name=object_name, source_type='knowledgebase', created_by=user_id, location=object_name, type=file_extension, status=5, size=size, parent_id=file.id ) logger.info(f"Created File with id: {fileDB.id}") File2document.objects.create( file_id=fileDB.id, document_id=documentKbm.id ) logger.info(f"Created File2document relation") child_count = File.objects.filter(parent_id=file.id).exclude(status=4).count() knowledgebase.doc_num = child_count knowledgebase.save() logger.info(f"Updated knowledgebase doc_num to {child_count}") except Exception as e: logger.error(f"Error in saveDocumentKbm: {str(e)}", exc_info=True) raise logger.info("saveDocumentKbm completed successfully") # 根据名称获取地址 @staticmethod def nameGetUrl(request): object_name = request.POST.get('object_name') bucket_name = request.POST.get('bucket_name') return MinioService.geturl(object_name,bucket_name) @staticmethod def geturl(name,bucket_name): """获取文件地址""" object_name = name if not object_name: return fail('Object name is required') try: # 生成一个预签名 URL,有效期为1小时 url = minio_client.presigned_get_object( bucket_name, object_name, expires=timedelta(hours=1) ) return success({'url': url}) except Exception as e: return fail(str(e)) @staticmethod def deleteFile(request): """根据名称删除文件""" bucket_name = request.POST.get('bucket_name') object_name = request.POST.get('object_name') if not object_name: return fail('需要删除的文件未找到') try: minio_client.remove_object(bucket_name, object_name) return success('删除成功') except S3Error as e: return fail(f'Failed to delete file: {str(e)}') except Exception as e: return fail(f'An unexpected error occurred: {str(e)}') @staticmethod def renameFile(request): """重命名 bucket 中的文件""" bucket_name = request.POST.get('bucket_name') old_name = request.POST.get('object_name') new_name = request.POST.get('new_name') if not old_name or not new_name: return fail('原文件名或新文件名未提供') # 检查新文件名是否包含后缀,如果没有则添加原文件的后缀 if '.' not in new_name: old_extension = old_name.split('.')[-1] if '.' in old_name else '' new_name = f"{new_name}.{old_extension}" if old_extension else new_name try: # 复制对象到新名称 result = minio_client.copy_object( bucket_name, new_name, CopySource(bucket_name, old_name) ) # 如果复制成功,删除原对象 if result: minio_client.remove_object(bucket_name, old_name) return success("修改昵称成功") else: return fail('文件重命名失败') except S3Error as e: return fail(f'重命名文件失败: {str(e)}') except Exception as e: return fail(f'发生意外错误: {str(e)}') @staticmethod def delete_bucket(request): """删除 bucket""" bucket_name = request.POST.get('bucket_name') if not bucket_name: return fail({'bucket_name 为空'}) try: # 开始事务 with transaction.atomic(): # 1. 从 MinIO 删除 bucket try: # 首先删除 bucket 中的所有对象 objects = minio_client.list_objects(bucket_name, recursive=True) for obj in objects: minio_client.remove_object(bucket_name, obj.object_name) # 然后删除 bucket minio_client.remove_bucket(bucket_name) except S3Error as e: return fail(f'MinIO 错误: {str(e)}') # 2. 从数据库中删除相应记录 try: bucket = Knowledgebase.objects.get(name=bucket_name) bucket.delete() except Knowledgebase.DoesNotExist: return fail('数据库中不存在该 bucket 记录') return fail(f'Bucket "{bucket_name}" 已成功删除') except Exception as e: return fail(f'删除 bucket 时发生错误: {str(e)}') @staticmethod def deleteBucket(request): return MinioService.delete_bucket(request) #切片 @staticmethod def readPdfSlice(request): """读取 PDF 切片,包括文本和图片""" bucket_name = request.POST.get('bucket_name') object_name = request.POST.get('object_name') start_page = int(request.POST.get('start_page', 1)) end_page = int(request.POST.get('end_page', -1)) if not bucket_name or not object_name: return fail('bucket_name 或 object_name 为空') try: # 从 MinIO 获取 PDF 文件 response = minio_client.get_object(bucket_name, object_name) pdf_content = BytesIO(response.read()) # 使用 PyMuPDF 读取 PDF doc = fitz.open(stream=pdf_content, filetype="pdf") total_pages = len(doc) # 调整页面范围 start_page = max(1, start_page) - 1 # 转换为从 0 开始的索引 end_page = min(total_pages, end_page if end_page > 0 else total_pages) # 读取指定页面范围的内容 result = [] for page_num in range(start_page, end_page): page = doc[page_num] text = page.get_text() # 提取图片 images = [] for img in page.get_images(): xref = img[0] base_image = doc.extract_image(xref) image_data = base_image["image"] image_format = base_image["ext"] image_base64 = base64.b64encode(image_data).decode('utf-8') images.append({ 'format': image_format, 'data': image_base64 }) result.append({ 'page_number': page_num + 1, 'content': text, 'images': images }) # 关闭连接 doc.close() pdf_content.close() response.close() info = { 'total_pages': total_pages, 'sliced_content': result } return success(info) except Exception as e: return fail(str(e))