Team Ai
Apppublic

Underground-Digital/Workflow-Engine

sourceHugging Faceupdated 2y agoView on Hugging Face
0likes
commands.py655 linesDownload Raw Back to api
1import base642import json3import logging4import secrets5from typing import Optional6 7import click8from flask import current_app9from werkzeug.exceptions import NotFound10 11from configs import dify_config12from constants.languages import languages13from core.rag.datasource.vdb.vector_factory import Vector14from core.rag.datasource.vdb.vector_type import VectorType15from core.rag.models.document import Document16from events.app_event import app_was_created17from extensions.ext_database import db18from extensions.ext_redis import redis_client19from libs.helper import email as email_validate20from libs.password import hash_password, password_pattern, valid_password21from libs.rsa import generate_key_pair22from models import Tenant23from models.dataset import Dataset, DatasetCollectionBinding, DocumentSegment24from models.dataset import Document as DatasetDocument25from models.model import Account, App, AppAnnotationSetting, AppMode, Conversation, MessageAnnotation26from models.provider import Provider, ProviderModel27from services.account_service import RegisterService, TenantService28 29 30@click.command("reset-password", help="Reset the account password.")31@click.option("--email", prompt=True, help="Account email to reset password for")32@click.option("--new-password", prompt=True, help="New password")33@click.option("--password-confirm", prompt=True, help="Confirm new password")34def reset_password(email, new_password, password_confirm):35    """36    Reset password of owner account37    Only available in SELF_HOSTED mode38    """39    if str(new_password).strip() != str(password_confirm).strip():40        click.echo(click.style("Passwords do not match.", fg="red"))41        return42 43    account = db.session.query(Account).filter(Account.email == email).one_or_none()44 45    if not account:46        click.echo(click.style("Account not found for email: {}".format(email), fg="red"))47        return48 49    try:50        valid_password(new_password)51    except:52        click.echo(click.style("Invalid password. Must match {}".format(password_pattern), fg="red"))53        return54 55    # generate password salt56    salt = secrets.token_bytes(16)57    base64_salt = base64.b64encode(salt).decode()58 59    # encrypt password with salt60    password_hashed = hash_password(new_password, salt)61    base64_password_hashed = base64.b64encode(password_hashed).decode()62    account.password = base64_password_hashed63    account.password_salt = base64_salt64    db.session.commit()65    click.echo(click.style("Password reset successfully.", fg="green"))66 67 68@click.command("reset-email", help="Reset the account email.")69@click.option("--email", prompt=True, help="Current account email")70@click.option("--new-email", prompt=True, help="New email")71@click.option("--email-confirm", prompt=True, help="Confirm new email")72def reset_email(email, new_email, email_confirm):73    """74    Replace account email75    :return:76    """77    if str(new_email).strip() != str(email_confirm).strip():78        click.echo(click.style("New emails do not match.", fg="red"))79        return80 81    account = db.session.query(Account).filter(Account.email == email).one_or_none()82 83    if not account:84        click.echo(click.style("Account not found for email: {}".format(email), fg="red"))85        return86 87    try:88        email_validate(new_email)89    except:90        click.echo(click.style("Invalid email: {}".format(new_email), fg="red"))91        return92 93    account.email = new_email94    db.session.commit()95    click.echo(click.style("Email updated successfully.", fg="green"))96 97 98@click.command(99    "reset-encrypt-key-pair",100    help="Reset the asymmetric key pair of workspace for encrypt LLM credentials. "101    "After the reset, all LLM credentials will become invalid, "102    "requiring re-entry."103    "Only support SELF_HOSTED mode.",104)105@click.confirmation_option(106    prompt=click.style(107        "Are you sure you want to reset encrypt key pair? This operation cannot be rolled back!", fg="red"108    )109)110def reset_encrypt_key_pair():111    """112    Reset the encrypted key pair of workspace for encrypt LLM credentials.113    After the reset, all LLM credentials will become invalid, requiring re-entry.114    Only support SELF_HOSTED mode.115    """116    if dify_config.EDITION != "SELF_HOSTED":117        click.echo(click.style("This command is only for SELF_HOSTED installations.", fg="red"))118        return119 120    tenants = db.session.query(Tenant).all()121    for tenant in tenants:122        if not tenant:123            click.echo(click.style("No workspaces found. Run /install first.", fg="red"))124            return125 126        tenant.encrypt_public_key = generate_key_pair(tenant.id)127 128        db.session.query(Provider).filter(Provider.provider_type == "custom", Provider.tenant_id == tenant.id).delete()129        db.session.query(ProviderModel).filter(ProviderModel.tenant_id == tenant.id).delete()130        db.session.commit()131 132        click.echo(133            click.style(134                "Congratulations! The asymmetric key pair of workspace {} has been reset.".format(tenant.id),135                fg="green",136            )137        )138 139 140@click.command("vdb-migrate", help="Migrate vector db.")141@click.option("--scope", default="all", prompt=False, help="The scope of vector database to migrate, Default is All.")142def vdb_migrate(scope: str):143    if scope in {"knowledge", "all"}:144        migrate_knowledge_vector_database()145    if scope in {"annotation", "all"}:146        migrate_annotation_vector_database()147 148 149def migrate_annotation_vector_database():150    """151    Migrate annotation datas to target vector database .152    """153    click.echo(click.style("Starting annotation data migration.", fg="green"))154    create_count = 0155    skipped_count = 0156    total_count = 0157    page = 1158    while True:159        try:160            # get apps info161            apps = (162                db.session.query(App)163                .filter(App.status == "normal")164                .order_by(App.created_at.desc())165                .paginate(page=page, per_page=50)166            )167        except NotFound:168            break169 170        page += 1171        for app in apps:172            total_count = total_count + 1173            click.echo(174                f"Processing the {total_count} app {app.id}. " + f"{create_count} created, {skipped_count} skipped."175            )176            try:177                click.echo("Creating app annotation index: {}".format(app.id))178                app_annotation_setting = (179                    db.session.query(AppAnnotationSetting).filter(AppAnnotationSetting.app_id == app.id).first()180                )181 182                if not app_annotation_setting:183                    skipped_count = skipped_count + 1184                    click.echo("App annotation setting disabled: {}".format(app.id))185                    continue186                # get dataset_collection_binding info187                dataset_collection_binding = (188                    db.session.query(DatasetCollectionBinding)189                    .filter(DatasetCollectionBinding.id == app_annotation_setting.collection_binding_id)190                    .first()191                )192                if not dataset_collection_binding:193                    click.echo("App annotation collection binding not found: {}".format(app.id))194                    continue195                annotations = db.session.query(MessageAnnotation).filter(MessageAnnotation.app_id == app.id).all()196                dataset = Dataset(197                    id=app.id,198                    tenant_id=app.tenant_id,199                    indexing_technique="high_quality",200                    embedding_model_provider=dataset_collection_binding.provider_name,201                    embedding_model=dataset_collection_binding.model_name,202                    collection_binding_id=dataset_collection_binding.id,203                )204                documents = []205                if annotations:206                    for annotation in annotations:207                        document = Document(208                            page_content=annotation.question,209                            metadata={"annotation_id": annotation.id, "app_id": app.id, "doc_id": annotation.id},210                        )211                        documents.append(document)212 213                vector = Vector(dataset, attributes=["doc_id", "annotation_id", "app_id"])214                click.echo(f"Migrating annotations for app: {app.id}.")215 216                try:217                    vector.delete()218                    click.echo(click.style(f"Deleted vector index for app {app.id}.", fg="green"))219                except Exception as e:220                    click.echo(click.style(f"Failed to delete vector index for app {app.id}.", fg="red"))221                    raise e222                if documents:223                    try:224                        click.echo(225                            click.style(226                                f"Creating vector index with {len(documents)} annotations for app {app.id}.",227                                fg="green",228                            )229                        )230                        vector.create(documents)231                        click.echo(click.style(f"Created vector index for app {app.id}.", fg="green"))232                    except Exception as e:233                        click.echo(click.style(f"Failed to created vector index for app {app.id}.", fg="red"))234                        raise e235                click.echo(f"Successfully migrated app annotation {app.id}.")236                create_count += 1237            except Exception as e:238                click.echo(239                    click.style(240                        "Error creating app annotation index: {} {}".format(e.__class__.__name__, str(e)), fg="red"241                    )242                )243                continue244 245    click.echo(246        click.style(247            f"Migration complete. Created {create_count} app annotation indexes. Skipped {skipped_count} apps.",248            fg="green",249        )250    )251 252 253def migrate_knowledge_vector_database():254    """255    Migrate vector database datas to target vector database .256    """257    click.echo(click.style("Starting vector database migration.", fg="green"))258    create_count = 0259    skipped_count = 0260    total_count = 0261    vector_type = dify_config.VECTOR_STORE262    upper_colletion_vector_types = {263        VectorType.MILVUS,264        VectorType.PGVECTOR,265        VectorType.RELYT,266        VectorType.WEAVIATE,267        VectorType.ORACLE,268        VectorType.ELASTICSEARCH,269    }270    lower_colletion_vector_types = {271        VectorType.ANALYTICDB,272        VectorType.CHROMA,273        VectorType.MYSCALE,274        VectorType.PGVECTO_RS,275        VectorType.TIDB_VECTOR,276        VectorType.OPENSEARCH,277        VectorType.TENCENT,278        VectorType.BAIDU,279        VectorType.VIKINGDB,280        VectorType.UPSTASH,281        VectorType.COUCHBASE,282        VectorType.OCEANBASE,283    }284    page = 1285    while True:286        try:287            datasets = (288                db.session.query(Dataset)289                .filter(Dataset.indexing_technique == "high_quality")290                .order_by(Dataset.created_at.desc())291                .paginate(page=page, per_page=50)292            )293        except NotFound:294            break295 296        page += 1297        for dataset in datasets:298            total_count = total_count + 1299            click.echo(300                f"Processing the {total_count} dataset {dataset.id}. {create_count} created, {skipped_count} skipped."301            )302            try:303                click.echo("Creating dataset vector database index: {}".format(dataset.id))304                if dataset.index_struct_dict:305                    if dataset.index_struct_dict["type"] == vector_type:306                        skipped_count = skipped_count + 1307                        continue308                collection_name = ""309                dataset_id = dataset.id310                if vector_type in upper_colletion_vector_types:311                    collection_name = Dataset.gen_collection_name_by_id(dataset_id)312                elif vector_type == VectorType.QDRANT:313                    if dataset.collection_binding_id:314                        dataset_collection_binding = (315                            db.session.query(DatasetCollectionBinding)316                            .filter(DatasetCollectionBinding.id == dataset.collection_binding_id)317                            .one_or_none()318                        )319                        if dataset_collection_binding:320                            collection_name = dataset_collection_binding.collection_name321                        else:322                            raise ValueError("Dataset Collection Binding not found")323                    else:324                        collection_name = Dataset.gen_collection_name_by_id(dataset_id)325 326                elif vector_type in lower_colletion_vector_types:327                    collection_name = Dataset.gen_collection_name_by_id(dataset_id).lower()328                else:329                    raise ValueError(f"Vector store {vector_type} is not supported.")330 331                index_struct_dict = {"type": vector_type, "vector_store": {"class_prefix": collection_name}}332                dataset.index_struct = json.dumps(index_struct_dict)333                vector = Vector(dataset)334                click.echo(f"Migrating dataset {dataset.id}.")335 336                try:337                    vector.delete()338                    click.echo(339                        click.style(f"Deleted vector index {collection_name} for dataset {dataset.id}.", fg="green")340                    )341                except Exception as e:342                    click.echo(343                        click.style(344                            f"Failed to delete vector index {collection_name} for dataset {dataset.id}.", fg="red"345                        )346                    )347                    raise e348 349                dataset_documents = (350                    db.session.query(DatasetDocument)351                    .filter(352                        DatasetDocument.dataset_id == dataset.id,353                        DatasetDocument.indexing_status == "completed",354                        DatasetDocument.enabled == True,355                        DatasetDocument.archived == False,356                    )357                    .all()358                )359 360                documents = []361                segments_count = 0362                for dataset_document in dataset_documents:363                    segments = (364                        db.session.query(DocumentSegment)365                        .filter(366                            DocumentSegment.document_id == dataset_document.id,367                            DocumentSegment.status == "completed",368                            DocumentSegment.enabled == True,369                        )370                        .all()371                    )372 373                    for segment in segments:374                        document = Document(375                            page_content=segment.content,376                            metadata={377                                "doc_id": segment.index_node_id,378                                "doc_hash": segment.index_node_hash,379                                "document_id": segment.document_id,380                                "dataset_id": segment.dataset_id,381                            },382                        )383 384                        documents.append(document)385                        segments_count = segments_count + 1386 387                if documents:388                    try:389                        click.echo(390                            click.style(391                                f"Creating vector index with {len(documents)} documents of {segments_count}"392                                f" segments for dataset {dataset.id}.",393                                fg="green",394                            )395                        )396                        vector.create(documents)397                        click.echo(click.style(f"Created vector index for dataset {dataset.id}.", fg="green"))398                    except Exception as e:399                        click.echo(click.style(f"Failed to created vector index for dataset {dataset.id}.", fg="red"))400                        raise e401                db.session.add(dataset)402                db.session.commit()403                click.echo(f"Successfully migrated dataset {dataset.id}.")404                create_count += 1405            except Exception as e:406                db.session.rollback()407                click.echo(408                    click.style("Error creating dataset index: {} {}".format(e.__class__.__name__, str(e)), fg="red")409                )410                continue411 412    click.echo(413        click.style(414            f"Migration complete. Created {create_count} dataset indexes. Skipped {skipped_count} datasets.", fg="green"415        )416    )417 418 419@click.command("convert-to-agent-apps", help="Convert Agent Assistant to Agent App.")420def convert_to_agent_apps():421    """422    Convert Agent Assistant to Agent App.423    """424    click.echo(click.style("Starting convert to agent apps.", fg="green"))425 426    proceeded_app_ids = []427 428    while True:429        # fetch first 1000 apps430        sql_query = """SELECT a.id AS id FROM apps a431            INNER JOIN app_model_configs am ON a.app_model_config_id=am.id432            WHERE a.mode = 'chat'433            AND am.agent_mode is not null434            AND (435				am.agent_mode like '%"strategy": "function_call"%'436                OR am.agent_mode  like '%"strategy": "react"%'437			)438            AND (439				am.agent_mode like '{"enabled": true%'440                OR am.agent_mode like '{"max_iteration": %'441			) ORDER BY a.created_at DESC LIMIT 1000442        """443 444        with db.engine.begin() as conn:445            rs = conn.execute(db.text(sql_query))446 447            apps = []448            for i in rs:449                app_id = str(i.id)450                if app_id not in proceeded_app_ids:451                    proceeded_app_ids.append(app_id)452                    app = db.session.query(App).filter(App.id == app_id).first()453                    apps.append(app)454 455            if len(apps) == 0:456                break457 458        for app in apps:459            click.echo("Converting app: {}".format(app.id))460 461            try:462                app.mode = AppMode.AGENT_CHAT.value463                db.session.commit()464 465                # update conversation mode to agent466                db.session.query(Conversation).filter(Conversation.app_id == app.id).update(467                    {Conversation.mode: AppMode.AGENT_CHAT.value}468                )469 470                db.session.commit()471                click.echo(click.style("Converted app: {}".format(app.id), fg="green"))472            except Exception as e:473                click.echo(click.style("Convert app error: {} {}".format(e.__class__.__name__, str(e)), fg="red"))474 475    click.echo(click.style("Conversion complete. Converted {} agent apps.".format(len(proceeded_app_ids)), fg="green"))476 477 478@click.command("add-qdrant-doc-id-index", help="Add Qdrant doc_id index.")479@click.option("--field", default="metadata.doc_id", prompt=False, help="Index field , default is metadata.doc_id.")480def add_qdrant_doc_id_index(field: str):481    click.echo(click.style("Starting Qdrant doc_id index creation.", fg="green"))482    vector_type = dify_config.VECTOR_STORE483    if vector_type != "qdrant":484        click.echo(click.style("This command only supports Qdrant vector store.", fg="red"))485        return486    create_count = 0487 488    try:489        bindings = db.session.query(DatasetCollectionBinding).all()490        if not bindings:491            click.echo(click.style("No dataset collection bindings found.", fg="red"))492            return493        import qdrant_client494        from qdrant_client.http.exceptions import UnexpectedResponse495        from qdrant_client.http.models import PayloadSchemaType496 497        from core.rag.datasource.vdb.qdrant.qdrant_vector import QdrantConfig498 499        for binding in bindings:500            if dify_config.QDRANT_URL is None:501                raise ValueError("Qdrant URL is required.")502            qdrant_config = QdrantConfig(503                endpoint=dify_config.QDRANT_URL,504                api_key=dify_config.QDRANT_API_KEY,505                root_path=current_app.root_path,506                timeout=dify_config.QDRANT_CLIENT_TIMEOUT,507                grpc_port=dify_config.QDRANT_GRPC_PORT,508                prefer_grpc=dify_config.QDRANT_GRPC_ENABLED,509            )510            try:511                client = qdrant_client.QdrantClient(**qdrant_config.to_qdrant_params())512                # create payload index513                client.create_payload_index(binding.collection_name, field, field_schema=PayloadSchemaType.KEYWORD)514                create_count += 1515            except UnexpectedResponse as e:516                # Collection does not exist, so return517                if e.status_code == 404:518                    click.echo(click.style(f"Collection not found: {binding.collection_name}.", fg="red"))519                    continue520                # Some other error occurred, so re-raise the exception521                else:522                    click.echo(523                        click.style(524                            f"Failed to create Qdrant index for collection: {binding.collection_name}.", fg="red"525                        )526                    )527 528    except Exception as e:529        click.echo(click.style("Failed to create Qdrant client.", fg="red"))530 531    click.echo(click.style(f"Index creation complete. Created {create_count} collection indexes.", fg="green"))532 533 534@click.command("create-tenant", help="Create account and tenant.")535@click.option("--email", prompt=True, help="Tenant account email.")536@click.option("--name", prompt=True, help="Workspace name.")537@click.option("--language", prompt=True, help="Account language, default: en-US.")538def create_tenant(email: str, language: Optional[str] = None, name: Optional[str] = None):539    """540    Create tenant account541    """542    if not email:543        click.echo(click.style("Email is required.", fg="red"))544        return545 546    # Create account547    email = email.strip()548 549    if "@" not in email:550        click.echo(click.style("Invalid email address.", fg="red"))551        return552 553    account_name = email.split("@")[0]554 555    if language not in languages:556        language = "en-US"557 558    name = name.strip()559 560    # generate random password561    new_password = secrets.token_urlsafe(16)562 563    # register account564    account = RegisterService.register(email=email, name=account_name, password=new_password, language=language)565 566    TenantService.create_owner_tenant_if_not_exist(account, name)567 568    click.echo(569        click.style(570            "Account and tenant created.\nAccount: {}\nPassword: {}".format(email, new_password),571            fg="green",572        )573    )574 575 576@click.command("upgrade-db", help="Upgrade the database")577def upgrade_db():578    click.echo("Preparing database migration...")579    lock = redis_client.lock(name="db_upgrade_lock", timeout=60)580    if lock.acquire(blocking=False):581        try:582            click.echo(click.style("Starting database migration.", fg="green"))583 584            # run db migration585            import flask_migrate586 587            flask_migrate.upgrade()588 589            click.echo(click.style("Database migration successful!", fg="green"))590 591        except Exception as e:592            logging.exception(f"Database migration failed: {e}")593        finally:594            lock.release()595    else:596        click.echo("Database migration skipped")597 598 599@click.command("fix-app-site-missing", help="Fix app related site missing issue.")600def fix_app_site_missing():601    """602    Fix app related site missing issue.603    """604    click.echo(click.style("Starting fix for missing app-related sites.", fg="green"))605 606    failed_app_ids = []607    while True:608        sql = """select apps.id as id from apps left join sites on sites.app_id=apps.id609where sites.id is null limit 1000"""610        with db.engine.begin() as conn:611            rs = conn.execute(db.text(sql))612 613            processed_count = 0614            for i in rs:615                processed_count += 1616                app_id = str(i.id)617 618                if app_id in failed_app_ids:619                    continue620 621                try:622                    app = db.session.query(App).filter(App.id == app_id).first()623                    tenant = app.tenant624                    if tenant:625                        accounts = tenant.get_accounts()626                        if not accounts:627                            print("Fix failed for app {}".format(app.id))628                            continue629 630                        account = accounts[0]631                        print("Fixing missing site for app {}".format(app.id))632                        app_was_created.send(app, account=account)633                except Exception as e:634                    failed_app_ids.append(app_id)635                    click.echo(click.style("Failed to fix missing site for app {}".format(app_id), fg="red"))636                    logging.exception(f"Fix app related site missing issue failed, error: {e}")637                    continue638 639            if not processed_count:640                break641 642    click.echo(click.style("Fix for missing app-related sites completed successfully!", fg="green"))643 644 645def register_commands(app):646    app.cli.add_command(reset_password)647    app.cli.add_command(reset_email)648    app.cli.add_command(reset_encrypt_key_pair)649    app.cli.add_command(vdb_migrate)650    app.cli.add_command(convert_to_agent_apps)651    app.cli.add_command(add_qdrant_doc_id_index)652    app.cli.add_command(create_tenant)653    app.cli.add_command(upgrade_db)654    app.cli.add_command(fix_app_site_missing)655