Underground-Digital/Workflow-Engine
0
1import logging2from collections.abc import Generator3 4import boto35from botocore.client import Config6from botocore.exceptions import ClientError7 8from configs import dify_config9from extensions.storage.base_storage import BaseStorage10 11logger = logging.getLogger(__name__)12 13 14class AwsS3Storage(BaseStorage):15 """Implementation for Amazon Web Services S3 storage."""16 17 def __init__(self):18 super().__init__()19 self.bucket_name = dify_config.S3_BUCKET_NAME20 if dify_config.S3_USE_AWS_MANAGED_IAM:21 logger.info("Using AWS managed IAM role for S3")22 23 session = boto3.Session()24 region_name = dify_config.S3_REGION25 self.client = session.client(service_name="s3", region_name=region_name)26 else:27 logger.info("Using ak and sk for S3")28 29 self.client = boto3.client(30 "s3",31 aws_secret_access_key=dify_config.S3_SECRET_KEY,32 aws_access_key_id=dify_config.S3_ACCESS_KEY,33 endpoint_url=dify_config.S3_ENDPOINT,34 region_name=dify_config.S3_REGION,35 config=Config(s3={"addressing_style": dify_config.S3_ADDRESS_STYLE}),36 )37 # create bucket38 try:39 self.client.head_bucket(Bucket=self.bucket_name)40 except ClientError as e:41 # if bucket not exists, create it42 if e.response["Error"]["Code"] == "404":43 self.client.create_bucket(Bucket=self.bucket_name)44 # if bucket is not accessible, pass, maybe the bucket is existing but not accessible45 elif e.response["Error"]["Code"] == "403":46 pass47 else:48 # other error, raise exception49 raise50 51 def save(self, filename, data):52 self.client.put_object(Bucket=self.bucket_name, Key=filename, Body=data)53 54 def load_once(self, filename: str) -> bytes:55 try:56 data = self.client.get_object(Bucket=self.bucket_name, Key=filename)["Body"].read()57 except ClientError as ex:58 if ex.response["Error"]["Code"] == "NoSuchKey":59 raise FileNotFoundError("File not found")60 else:61 raise62 return data63 64 def load_stream(self, filename: str) -> Generator:65 try:66 response = self.client.get_object(Bucket=self.bucket_name, Key=filename)67 yield from response["Body"].iter_chunks()68 except ClientError as ex:69 if ex.response["Error"]["Code"] == "NoSuchKey":70 raise FileNotFoundError("File not found")71 else:72 raise73 74 def download(self, filename, target_filepath):75 self.client.download_file(self.bucket_name, filename, target_filepath)76 77 def exists(self, filename):78 try:79 self.client.head_object(Bucket=self.bucket_name, Key=filename)80 return True81 except:82 return False83 84 def delete(self, filename):85 self.client.delete_object(Bucket=self.bucket_name, Key=filename)86 