diff --git a/service/file_sync_service.py b/service/file_sync_service.py index c4bbb17..a341270 100644 --- a/service/file_sync_service.py +++ b/service/file_sync_service.py @@ -9,6 +9,45 @@ from pathlib import Path from .twoway_sync_folder_service import * import shutil +@server.PromptServer.instance.routes.post('/workspace/file/save') +async def save_file(request): + reqJson = await request.json() + data = await asyncio.to_thread(save_file_sync, reqJson) + return web.json_response(data, content_type='application/json') + +def save_file_sync(reqJson): + current_path = reqJson.get('path') + new_json_data_str = reqJson.get('json') + try: + new_json_data = json.loads(new_json_data_str) + except json.JSONDecodeError: + return { "error": "New JSON data is not valid JSON."} + expected_id = new_json_data.get('extra', {}).get('workspace_info', {}).get('id', None) + + try: + with open(current_path, 'r', encoding='utf-8') as file: + current_json_data = json.load(file) + # Extract the id from the current file's content + current_id = current_json_data.get('extra', {}).get('workspace_info', {}).get('id', None) + + # Compare the current id with the expected id + if current_id is not None and current_id != expected_id: + return {"error": "Mismatching workspace_info.id."} + except json.JSONDecodeError: + # If there's a JSON decode error, it means the file is not a valid JSON + return { "error": "Existing file is not a valid JSON."} + except FileNotFoundError: + # This block is optional since os.path.exists already checks for the file's existence + return {"error": "File not found."} + + # If checks pass, write the new JSON data to the file + try: + with open(current_path, 'w', encoding='utf-8') as file: + file.write(new_json_data_str) + return {"success": True} + except Exception as e: + return { "error": str(e)} + @server.PromptServer.instance.routes.post('/workspace/file/rename') async def rename_file(request): reqJson = await request.json() diff --git a/service/scan_my_workflows_folder.py b/service/scan_my_workflows_folder.py index 390214c..a130c82 100644 --- a/service/scan_my_workflows_folder.py +++ b/service/scan_my_workflows_folder.py @@ -1,75 +1,45 @@ -import os -import json -import uuid -import tempfile -import server -import asyncio from aiohttp import web +import asyncio +import json +import os +import glob +from threading import Lock +import server -id_cache = {} -# Scan .json and folders in the given path -@server.PromptServer.instance.routes.post("/workspace/scan_my_workflows_files") -async def scan_local_new_files(request): +@server.PromptServer.instance.routes.post('/workspace/file/scan_dup_id') +async def scan_my_workflows_folder_recursive(request): reqJson = await request.json() - path = reqJson['path'] - print("Scanning path: ", path) - - fileList = await asyncio.to_thread(folder_handle, path) - return web.Response(text=json.dumps(fileList), content_type='application/json') + path = reqJson.get('path') + data = await asyncio.to_thread(scan_folder_sync, path) + return web.json_response(data, content_type='application/json') -def folder_handle(path): - global id_cache - fileList = [] - for item in os.listdir(path): - try: - item_path = os.path.join(path, item) - if os.path.isfile(item_path) and item_path.endswith('.json'): - stats = os.stat(item_path) - creation_time = stats.st_ctime - update_time = stats.st_mtime - with open(item_path, 'r', encoding='utf-8') as f: - file_handle(item, f, fileList, item_path, creation_time, update_time) +def scan_folder_sync(path): + id_cache = {} # Reset cache for each call + id_cache_lock = Lock() # Ensure thread safety - elif os.path.isdir(item_path): - fileList.append({ - 'name': item, - 'type': 'folder', - }) - except Exception as e: - print(f"Error scanning file {item}: {e}") - return fileList + def scan_folder(path): + for json_file in glob.glob(os.path.join(path, '**/*.json'), recursive=True): + with open(json_file, 'r', encoding='utf-8') as file: + try: + data = json.load(file) + if 'extra' in data and 'workspace_info' in data['extra'] and 'id' in data['extra']['workspace_info']: + workspace_id = data['extra']['workspace_info']['id'] + file_info = { + 'path': json_file, + 'createTime': os.path.getctime(json_file) + } + + with id_cache_lock: + if workspace_id in id_cache: + id_cache[workspace_id].append(file_info) + else: + id_cache[workspace_id] = [file_info] + except json.JSONDecodeError: + print(f"Error decoding JSON from file: {json_file}") + + scan_folder(path) + + duplicates = {workspace_id: infos for workspace_id, infos in id_cache.items() if len(infos) > 1} + return {'duplicates': duplicates} -def file_handle(name, file, fileList, file_path, creation_time, update_time): - json_data = json.load(file) - if 'extra' in json_data and 'workspace_info' in json_data['extra'] and 'id' in json_data['extra']['workspace_info']: - workflow_id = json_data['extra']['workspace_info']['id'] - if workflow_id in id_cache and id_cache[workflow_id] < creation_time: - # File with later creation time found, assign new ID - workflow_id = str(uuid.uuid4()) - json_data['extra']['workspace_info']['id'] = workflow_id - atomic_json_update(file_path, json_data) - id_cache[workflow_id] = creation_time - else: - # If ID does not exist, generate a new UUID and add it to the JSON data - workflow_id = str(uuid.uuid4()) - json_data.setdefault('extra', {}).setdefault('workspace_info', {})['id'] = workflow_id - atomic_json_update(file_path, json_data) - id_cache[workflow_id] = creation_time - fileInfo = { - 'json': json.dumps(json_data), - 'name': name, - 'type': "workflow", - 'id': workflow_id, - 'createTime': creation_time, - 'updateTime': update_time - } - fileList.append(fileInfo) -def atomic_json_update(filepath, data): - # Generate a temporary file - dir_name, file_name = os.path.split(filepath) - with tempfile.NamedTemporaryFile(mode='w', encoding='utf-8', dir=dir_name, delete=False) as tmp_file: - json.dump(data, tmp_file, indent=4) - temp_name = tmp_file.name - # Replace the old file with the new file atomically - os.replace(temp_name, filepath) diff --git a/ui/src/App.tsx b/ui/src/App.tsx index 866e8d3..ed4f604 100644 --- a/ui/src/App.tsx +++ b/ui/src/App.tsx @@ -73,7 +73,6 @@ export default function App() { const graphJson = JSON.stringify(app.graph.serialize()); await Promise.all([ workflowsTable?.updateFlow(curFlowID.current, { - lastSavedJson: graphJson, json: graphJson, }), changelogsTable?.create({ @@ -176,7 +175,7 @@ export default function App() { const checkIsDirty = async () => { if (curFlowID.current != null) { - const curWorkflow = await workflowsTable?.get(curFlowID.current, false); + const curWorkflow = await workflowsTable?.get(curFlowID.current); return !!curWorkflow && checkIsDirtyImpl(curWorkflow); } return false; diff --git a/ui/src/RecentFilesDrawer/FilesListFolderItem.tsx b/ui/src/RecentFilesDrawer/FilesListFolderItem.tsx index c1c683d..ef0fc80 100644 --- a/ui/src/RecentFilesDrawer/FilesListFolderItem.tsx +++ b/ui/src/RecentFilesDrawer/FilesListFolderItem.tsx @@ -67,7 +67,7 @@ export default memo(function FilesListFolderItem({ folder }: Props) { parentFolderID: folder.id, }); } else if (!isFolder(draggingFile)) { - await workflowsTable?.updateFlow(draggingFile.id, { + await workflowsTable?.updateFolder(draggingFile.id, { parentFolderID: folder.id, }); } diff --git a/ui/src/RecentFilesDrawer/ItemsList.tsx b/ui/src/RecentFilesDrawer/ItemsList.tsx index 3f4feb8..65fbfb4 100644 --- a/ui/src/RecentFilesDrawer/ItemsList.tsx +++ b/ui/src/RecentFilesDrawer/ItemsList.tsx @@ -35,7 +35,7 @@ export default function ItemsList({ parentFolderID: parentFolderID, }); } else { - await workflowsTable?.updateFlow(draggingFile.id, { + await workflowsTable?.updateFolder(draggingFile.id, { parentFolderID: parentFolderID, }); } diff --git a/ui/src/apis/TwowaySyncApi.ts b/ui/src/apis/TwowaySyncApi.ts index 6894cf2..24b4d84 100644 --- a/ui/src/apis/TwowaySyncApi.ts +++ b/ui/src/apis/TwowaySyncApi.ts @@ -4,6 +4,7 @@ import { userSettingsTable } from "../db-tables/WorkspaceDB"; import { indexdb } from "../db-tables/indexdb"; import { Workflow } from "../types/dbTypes"; import { genAbsPathByRelPath, sanitizeAbsPath } from "../utils/OsPathUtils"; +import { TwowayFolderSyncAPI } from "./TwowaySyncFolderApi"; export namespace TwowaySyncAPI { async function genWorkflowAbsPath({ @@ -99,16 +100,57 @@ export namespace TwowaySyncAPI { } } + interface DuplicatesResponse { + duplicates: { + [id: string]: Array<{ + path: string; + createTime: number; + }>; + }; + } + export async function scanMyWorkflowsDupId() { + const myWorkflowsDir = + await userSettingsTable?.getSetting("myWorkflowsDir"); + const absPath = sanitizeAbsPath(myWorkflowsDir!); + try { + const response = await fetch("/workspace/file/scan_dup_id", { + method: "POST", + headers: { + "Content-Type": "application/json", + }, + body: JSON.stringify({ + path: absPath, + }), + }); + const result = (await response.json()) as DuplicatesResponse; + console.log("scanMyWorkflowsDupId", result); + return true; + } catch (error) { + console.error("Error deleting file:", error); + } + } + export async function saveWorkflow(workflow: Workflow) { console.log("🥳saveWorkflow", workflow); - const file_path = await genWorkflowAbsPath(workflow); + const absPath = await genWorkflowAbsPath(workflow); const json = workflow.json; const flow = JSON.parse(json); flow.extra[COMFYSPACE_TRACKING_FIELD_NAME] = { id: workflow.id, }; - await updateFile(file_path, JSON.stringify(flow)); + const response = await fetch("/workspace/file/save", { + method: "POST", + headers: { + "Content-Type": "application/json", + }, + body: JSON.stringify({ + path: absPath, + json: json, + }), + }); + const result = await response.text(); + return result; } export async function deleteWorkflow(workflow: Workflow) { const absPath = await genWorkflowAbsPath(workflow); @@ -234,6 +276,8 @@ export type ScanLocalFile = { name: string; id: string; json: string; + createTime: number; + updateTime: number; }; export type ScanLocalFolder = { type: "folder"; diff --git a/ui/src/components/EditFlowName.tsx b/ui/src/components/EditFlowName.tsx index aa025d0..5033b3f 100644 --- a/ui/src/components/EditFlowName.tsx +++ b/ui/src/components/EditFlowName.tsx @@ -61,7 +61,7 @@ export default function EditFlowName({ "The name is duplicated, please modify it and submit again.", ); } else { - await workflowsTable?.updateFlow(curFlowID, { + await workflowsTable?.updateName(curFlowID, { name: trimEditName, }); updateFlowName(trimEditName); diff --git a/ui/src/components/VersionHistoryDrawer.tsx b/ui/src/components/VersionHistoryDrawer.tsx index d65e1c4..4c2b16a 100644 --- a/ui/src/components/VersionHistoryDrawer.tsx +++ b/ui/src/components/VersionHistoryDrawer.tsx @@ -147,7 +147,6 @@ export function VersionHistoryDrawer({ onClose }: { onClose: () => void }) { } app.loadGraphData(JSON.parse(version.json)); workflowsTable?.updateFlow(curFlowID!, { - lastSavedJson: version.json, json: version.json, }); toast({ @@ -191,7 +190,7 @@ export function VersionHistoryDrawer({ onClose }: { onClose: () => void }) { } app.loadGraphData(JSON.parse(c.json)); workflowsTable?.updateFlow(curFlowID!, { - lastSavedJson: c.json, + json: c.json, }); onClose(); }} diff --git a/ui/src/db-tables/DiskFileUtils.ts b/ui/src/db-tables/DiskFileUtils.ts index 899720f..69c9a9c 100644 --- a/ui/src/db-tables/DiskFileUtils.ts +++ b/ui/src/db-tables/DiskFileUtils.ts @@ -9,7 +9,7 @@ import { Workflow } from "../types/dbTypes"; import { TwowayFolderSyncAPI } from "../apis/TwowaySyncFolderApi"; export async function saveJsonFileMyWorkflows(workflow: Workflow) { - console.log("saveJsonFileMyWorkflows", workflow); + console.warn("saveJsonFileMyWorkflows", workflow); const file_path = await generateFilePath(workflow); if (file_path == null) { return; diff --git a/ui/src/db-tables/FoldersTable.ts b/ui/src/db-tables/FoldersTable.ts index c2ebba0..2814218 100644 --- a/ui/src/db-tables/FoldersTable.ts +++ b/ui/src/db-tables/FoldersTable.ts @@ -112,7 +112,7 @@ export class FoldersTable extends TableBase { await workflowsTable?.deleteFlow(flow.id); break; case EFlowOperationType.MOVE_TO_ROOT_FOLDER: - await workflowsTable?.updateFlow(flow.id, { + await workflowsTable?.updateFolder(flow.id, { parentFolderID: undefined, }); break; diff --git a/ui/src/db-tables/MediaTable.ts b/ui/src/db-tables/MediaTable.ts index 5fdbd0a..d15cf18 100644 --- a/ui/src/db-tables/MediaTable.ts +++ b/ui/src/db-tables/MediaTable.ts @@ -34,7 +34,7 @@ export class MediaTable extends TableBase { //link media to workflow const workflow = await workflowsTable?.get(input.workflowID); const newMedia = new Set(workflow?.mediaIDs ?? []).add(md.id); - await workflowsTable?.updateFlow(input.workflowID, { + await workflowsTable?.updateMetaInfo(input.workflowID, { mediaIDs: Array.from(newMedia), coverMediaPath: md.localPath, }); diff --git a/ui/src/db-tables/WorkflowsTable.ts b/ui/src/db-tables/WorkflowsTable.ts index afaa44f..ccb76fb 100644 --- a/ui/src/db-tables/WorkflowsTable.ts +++ b/ui/src/db-tables/WorkflowsTable.ts @@ -131,79 +131,69 @@ export class WorkflowsTable extends TableBase { ): Promise { throw new Error("Method not allowed."); } + public async updateMetaInfo( id: string, change: Omit, "id" | "name" | "parentFolderID" | "json">, ): Promise { //update indexdb await indexdb.workflows.update(id, change); - const newWorkflow = (await this.get(id, false)) ?? null; + const newWorkflow = (await indexdb.workflows.get(id)) ?? null; //update curWorkflow RAM if (this._curWorkflow && this._curWorkflow.id === id) { this._curWorkflow = newWorkflow; } - await this.saveDiskDB(); + this.saveDiskDB(); return newWorkflow; } - public async updateFlow(id: string, input: Partial) { - // console.log("updateFlow", id, input); - const before = await this.get(id); - if (before == null) { + async updateFolder(id: string, change: Pick) { + const twoWaySyncEnabled = await userSettingsTable?.getSetting("twoWaySync"); + const oldWorkflow = await this.get(id); + const newWorkflow = await this.updateMetaInfo(id, change as any); + if (!newWorkflow || !oldWorkflow) { return; } - const after = { - ...before, - ...input, - id, - }; - const beforeStr = JSON.stringify(before); - const afterStr = JSON.stringify(after); - if (beforeStr === afterStr) { - // no change detected - return; - } - const newWorkflow: Workflow = after; - // When modifying the associated tag or modifying the directory, updateTime is not modified. - const updateKey = Object.keys(input); - const isModifyingTagOrFolder = - updateKey.length === 1 && - ["tags", "parentFolderID"].includes(updateKey[0]); - if (!isModifyingTagOrFolder) { - newWorkflow.updateTime = Date.now(); + if (twoWaySyncEnabled) { + await TwowaySyncAPI.moveWorkflow(oldWorkflow, change.parentFolderID!); + } else { + await saveJsonFileMyWorkflows(newWorkflow); + await deleteJsonFileMyWorkflows(oldWorkflow); // update parent folder updateTime - if (newWorkflow.parentFolderID != null) { - await foldersTable?.update(newWorkflow.parentFolderID, { + if (newWorkflow?.parentFolderID != null) { + await indexdb.folders?.update(newWorkflow.parentFolderID, { updateTime: Date.now(), }); } } - //update indexdb - await indexdb.workflows.update(id, newWorkflow); - //update curWorkflow RAM - if (this._curWorkflow && this._curWorkflow.id === id) { - this._curWorkflow = newWorkflow; + } + public async updateName(id: string, change: Pick) { + const before = await this.get(id); + const twoWaySyncEnabled = await userSettingsTable?.getSetting("twoWaySync"); + if (twoWaySyncEnabled) { + before && + (await TwowaySyncAPI.renameWorkflow(before, change.name + ".json")); } - await this.saveDiskDB(); - // save to my_workflows/ - if ("name" in input || "parentFolderID" in input) { - const twoWaySyncEnabled = - await userSettingsTable?.getSetting("twoWaySync"); - // renamed file or moved file folder - if (twoWaySyncEnabled) { - if ("name" in input) { - await TwowaySyncAPI.renameWorkflow(before, input["name"]! + ".json"); - } else if ("parentFolderID" in input) { - await TwowaySyncAPI.moveWorkflow(before, input["parentFolderID"]!); - } - } else { - await saveJsonFileMyWorkflows(after); - await deleteJsonFileMyWorkflows(before); - } + } + public async updateFlow(id: string, input: Pick) { + const before = await this.get(id); + if (before == null) { return; } - if (input.json != null) { - await saveJsonFileMyWorkflows(after); + const beforeStr = JSON.stringify(before.json); + const afterStr = JSON.stringify(input.json); + if (beforeStr === afterStr) { + // no change detected + return; + } + console.log("updateFlow", id, input); + const after = await this.updateMetaInfo(id, input as any); + // save to my_workflows/ + const twoWaySyncEnabled = await userSettingsTable?.getSetting("twoWaySync"); + if (twoWaySyncEnabled) { + after && TwowaySyncAPI.saveWorkflow(after); + } else { + after && (await saveJsonFileMyWorkflows(after)); } } @@ -239,7 +229,13 @@ export class WorkflowsTable extends TableBase { tags: [], }; newWorkflows.push(newWorkflow); - await saveJsonFileMyWorkflows(newWorkflow); + const twoWaySyncEnabled = + await userSettingsTable?.getSetting("twoWaySync"); + if (twoWaySyncEnabled) { + TwowaySyncAPI.creatWorkflow(newWorkflow); + } else { + await saveJsonFileMyWorkflows(newWorkflow); + } } await indexdb.workflows.bulkAdd(newWorkflows); this.saveDiskDB(); diff --git a/ui/src/gallery/components/GalleryMediaItem.tsx b/ui/src/gallery/components/GalleryMediaItem.tsx index 2f55b29..2a09f06 100644 --- a/ui/src/gallery/components/GalleryMediaItem.tsx +++ b/ui/src/gallery/components/GalleryMediaItem.tsx @@ -91,7 +91,7 @@ const GalleryMediaItem: React.FC = ({ icon={isCover ? : } aria-label="set as cover" onClick={() => { - workflowsTable?.updateFlow(curFlowID, { + workflowsTable?.updateMetaInfo(curFlowID, { coverMediaPath: media.localPath, }); setCoverPath(media.localPath); @@ -147,7 +147,7 @@ const GalleryMediaItem: React.FC = ({ if (!res) return; await mediaTable?.delete(media.id); if (isCover) { - await workflowsTable?.updateFlow(curFlowID, { + await workflowsTable?.updateMetaInfo(curFlowID, { coverMediaPath: undefined, }); setCoverPath(""); diff --git a/ui/src/share/ShareDialog.tsx b/ui/src/share/ShareDialog.tsx index 8fed848..b2b5d6f 100644 --- a/ui/src/share/ShareDialog.tsx +++ b/ui/src/share/ShareDialog.tsx @@ -73,7 +73,7 @@ export default function ShareDialog({ onClose }: Props) { cloudID && localID && - (await workflowsTable?.updateFlow(localID, { + (await workflowsTable?.updateMetaInfo(localID, { cloudID: cloudID, cloudURL: cloudHostRef.current + "/workflow/" + cloudID, }));