Underground-Digital/Workflow-Engine
0
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 