Home

mweinbach / agent-coworker

publicmweinbach/agent-coworker
Overview Code History Branches Pull requestsIssuesInsights
main
HomeOverview Code PRsIssues
mweinbach/agent-coworker/apps/desktop/electron/services/workspaceDirectoryWatcher.ts
Raw
1import { type FSWatcher, watch as watchFileSystem } from "node:fs";2import fs from "node:fs/promises";3import path from "node:path";4 5import {6  createWorkspaceFileChangeEvent,7
type
WorkspaceFileChangeEvent,
8 type WorkspaceFileChangeKind,
9} from "../../../../src/filesystem/workspaceFileEvents";
10import { isPathInside } from "../../../../src/utils/paths";
11 
12export type WorkspaceDirectoryWatchScope = {
13 workspaceId: string;
14 rootPath: string;
15};
16 
17export type DirectoryWatchListener = (event: WorkspaceFileChangeEvent) => void;
18 
19type WatchFactory = (
20 rootPath: string,
21 listener: (eventType: "rename" | "change", filename: string | Buffer | null) => void,
22) => Pick<FSWatcher, "close">;
23 
24export type WorkspaceDirectoryWatcherOptions = {
25 debounceMs?: number;
26 pathExists?: (candidatePath: string) => Promise<boolean>;
27 watch?: WatchFactory;
28};
29 
30type PendingWatchEvent = {
31 eventType: "rename" | "change";
32 path: string;
33};
34 
35type ActiveWatch = {
36 debounceTimer: ReturnType<typeof setTimeout> | null;
37 pendingByPath: Map<string, PendingWatchEvent>;
38 rootPath: string;
39 subscribers: Map<string, DirectoryWatchListener>;
40 watcher: Pick<FSWatcher, "close">;
41 workspaceId: string;
42};
43 
44const DEFAULT_WATCH_DEBOUNCE_MS = 40;
45 
46async function defaultPathExists(candidatePath: string): Promise<boolean> {
47 try {
48 await fs.lstat(candidatePath);
49 return true;
50 } catch {
51 return false;
52 }
53}
54 
55function defaultWatchFactory(
56 rootPath: string,
57 listener: Parameters<WatchFactory>[1],
58): Pick<FSWatcher, "close"> {
59 return watchFileSystem(rootPath, { recursive: true }, listener);
60}
61 
62function watchScopeKey(scope: WorkspaceDirectoryWatchScope): string {
63 return `${scope.workspaceId}\0${path.resolve(scope.rootPath)}`;
64}
65 
66export class WorkspaceDirectoryWatcher {
67 private readonly activeByScope = new Map<string, ActiveWatch>();
68 private readonly debounceMs: number;
69 private readonly pathExists: (candidatePath: string) => Promise<boolean>;
70 private readonly watchFactory: WatchFactory;
71 
72 constructor(options: WorkspaceDirectoryWatcherOptions = {}) {
73 this.debounceMs = options.debounceMs ?? DEFAULT_WATCH_DEBOUNCE_MS;
74 this.pathExists = options.pathExists ?? defaultPathExists;
75 this.watchFactory = options.watch ?? defaultWatchFactory;
76 }
77 
78 watch(
79 scope: WorkspaceDirectoryWatchScope,
80 subscriberId: string,
81 listener: DirectoryWatchListener,
82 ): boolean {
83 const key = watchScopeKey(scope);
84 const existing = this.activeByScope.get(key);
85 if (existing) {
86 existing.subscribers.set(subscriberId, listener);
87 return true;
88 }
89 
90 const rootPath = path.resolve(scope.rootPath);
91 let active: ActiveWatch | null = null;
92 try {
93 const watcher = this.watchFactory(rootPath, (eventType, filename) => {
94 if (active) {
95 this.queueRawEvent(active, eventType, filename);
96 }
97 });
98 active = {
99 debounceTimer: null,
100 pendingByPath: new Map(),
101 rootPath,
102 subscribers: new Map([[subscriberId, listener]]),
103 watcher,
104 workspaceId: scope.workspaceId,
105 };
106 } catch {
107 return false;
108 }
109 this.activeByScope.set(key, active);
110 return true;
111 }
112 
113 unwatch(scope: WorkspaceDirectoryWatchScope, subscriberId: string): void {
114 const key = watchScopeKey(scope);
115 const active = this.activeByScope.get(key);
116 if (!active) {
117 return;
118 }
119 active.subscribers.delete(subscriberId);
120 if (active.subscribers.size > 0) {
121 return;
122 }
123 this.closeWatch(key, active);
124 }
125 
126 unwatchSubscriber(subscriberId: string): void {
127 for (const [key, active] of this.activeByScope) {
128 active.subscribers.delete(subscriberId);
129 if (active.subscribers.size === 0) {
130 this.closeWatch(key, active);
131 }
132 }
133 }
134 
135 dispose(): void {
136 for (const [key, active] of this.activeByScope) {
137 this.closeWatch(key, active);
138 }
139 }
140 
141 private queueRawEvent(
142 active: ActiveWatch,
143 eventType: "rename" | "change",
144 filename: string | Buffer | null,
145 ): void {
146 const relativePath = filename?.toString().trim() ?? "";
147 const changedPath = relativePath
148 ? path.resolve(active.rootPath, relativePath)
149 : active.rootPath;
150 if (changedPath !== active.rootPath && !isPathInside(active.rootPath, changedPath)) {
151 return;
152 }
153 active.pendingByPath.set(changedPath, { eventType, path: changedPath });
154 if (active.debounceTimer) {
155 clearTimeout(active.debounceTimer);
156 }
157 active.debounceTimer = setTimeout(() => {
158 active.debounceTimer = null;
159 void this.flush(active);
160 }, this.debounceMs);
161 }
162 
163 private async flush(active: ActiveWatch): Promise<void> {
164 const pending = [...active.pendingByPath.values()];
165 active.pendingByPath.clear();
166 if (pending.length === 0 || active.subscribers.size === 0) {
167 return;
168 }
169 
170 const modifiedPaths = pending
171 .filter((event) => event.eventType === "change")
172 .map((event) => event.path);
173 if (modifiedPaths.length > 0) {
174 this.emit(active, "modify", modifiedPaths);
175 }
176 
177 const renameCandidates = pending.filter((event) => event.eventType === "rename");
178 if (renameCandidates.length === 0) {
179 return;
180 }
181 const existence = await Promise.all(
182 renameCandidates.map(async (event) => ({
183 exists: await this.pathExists(event.path),
184 path: event.path,
185 })),
186 );
187 if (active.subscribers.size === 0) {
188 return;
189 }
190 const addedPaths = existence.filter((entry) => entry.exists).map((entry) => entry.path);
191 const removedPaths = existence.filter((entry) => !entry.exists).map((entry) => entry.path);
192 
193 if (addedPaths.length > 0 && removedPaths.length > 0) {
194 this.emit(active, "rename", [...removedPaths, ...addedPaths]);
195 return;
196 }
197 if (addedPaths.length > 0) {
198 this.emit(active, "add", addedPaths);
199 }
200 if (removedPaths.length > 0) {
201 this.emit(active, "remove", removedPaths);
202 }
203 }
204 
205 private emit(active: ActiveWatch, kind: WorkspaceFileChangeKind, changedPaths: string[]): void {
206 const event = createWorkspaceFileChangeEvent({
207 workspaceId: active.workspaceId,
208 rootPath: active.rootPath,
209 kind,
210 changedPaths,
211 });
212 for (const listener of active.subscribers.values()) {
213 listener(event);
214 }
215 }
216 
217 private closeWatch(key: string, active: ActiveWatch): void {
218 if (active.debounceTimer) {
219 clearTimeout(active.debounceTimer);
220 }
221 active.watcher.close();
222 this.activeByScope.delete(key);
223 }
224}
225