splitted updateFlow into updateMetaInfo updateName updateFolder for better mainainability

This commit is contained in:
Weixuan Fu
2024-02-27 12:12:31 +08:00
parent 4fbf25c403
commit d0ae998f11
14 changed files with 181 additions and 134 deletions
+39
View File
@@ -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()
+38 -68
View File
@@ -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)
+1 -2
View File
@@ -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;
@@ -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,
});
}
+1 -1
View File
@@ -35,7 +35,7 @@ export default function ItemsList({
parentFolderID: parentFolderID,
});
} else {
await workflowsTable?.updateFlow(draggingFile.id, {
await workflowsTable?.updateFolder(draggingFile.id, {
parentFolderID: parentFolderID,
});
}
+46 -2
View File
@@ -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";
+1 -1
View File
@@ -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);
+1 -2
View File
@@ -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();
}}
+1 -1
View File
@@ -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;
+1 -1
View File
@@ -112,7 +112,7 @@ export class FoldersTable extends TableBase<Folder> {
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;
+1 -1
View File
@@ -34,7 +34,7 @@ export class MediaTable extends TableBase<Media> {
//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,
});
+47 -51
View File
@@ -131,79 +131,69 @@ export class WorkflowsTable extends TableBase<Workflow> {
): Promise<Workflow | null> {
throw new Error("Method not allowed.");
}
public async updateMetaInfo(
id: string,
change: Omit<Partial<Workflow>, "id" | "name" | "parentFolderID" | "json">,
): Promise<Workflow | null> {
//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<Workflow>) {
// console.log("updateFlow", id, input);
const before = await this.get(id);
if (before == null) {
async updateFolder(id: string, change: Pick<Workflow, "parentFolderID">) {
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<Workflow, "name">) {
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<Workflow, "json">) {
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<Workflow> {
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();
@@ -91,7 +91,7 @@ const GalleryMediaItem: React.FC<GalleryMediaItemProps> = ({
icon={isCover ? <IconPinFilled size={19} /> : <IconPin size={19} />}
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<GalleryMediaItemProps> = ({
if (!res) return;
await mediaTable?.delete(media.id);
if (isCover) {
await workflowsTable?.updateFlow(curFlowID, {
await workflowsTable?.updateMetaInfo(curFlowID, {
coverMediaPath: undefined,
});
setCoverPath("");
+1 -1
View File
@@ -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,
}));