Team Ai
Apppublic

Underground-Digital/Workflow-Engine

sourceHugging Faceupdated 2y agoView on Hugging Face
0likes
aws_s3_storage.py86 linesDownload Raw Back to storage
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