feat(FN-4748): complete Step 3 — attach tracking state service to all project stores

Fusion-Task-Id: FN-4748
Fusion-Task-Lineage: b540f985-c20f-438e-8290-a07531c6612c
This commit is contained in:
Fusion (runfusion.ai)
2026-05-16 09:57:05 -07:00
committed by gsxdsm
parent d7b855c31d
commit e9933cfa94
5 changed files with 123 additions and 28 deletions

View File

@@ -255,9 +255,10 @@ describe("GitHubTrackingStateService", () => {
expect(mockResolveGithubTrackingAuth).toHaveBeenCalledTimes(2);
});
it("closes issue for late-registered project stores after service start", async () => {
it("closes issue for late-registered project stores after attach", async () => {
const lateStore = new MockStore();
service.start();
service.attach(lateStore as unknown as TaskStore);
lateStore.emit("task:moved", { task: createTask({ id: "FN-late" }), from: "todo", to: "done" });
await flushAsync();

View File

@@ -14,6 +14,8 @@ import {
getOrCreateProjectStore,
evictProjectStore,
evictAllProjectStores,
listRegisteredProjectStores,
onProjectStoreRegistered,
setOnProjectFirstCreated,
} from "../project-store-resolver.js";
import type { TaskStore } from "@fusion/core";
@@ -79,6 +81,29 @@ describe("project-store-resolver", () => {
expect(createdStores[0].projectId).toBe("proj_abc");
});
it("lists registered stores", async () => {
const storeA = await getOrCreateProjectStore("proj_list_a");
const storeB = await getOrCreateProjectStore("proj_list_b");
const registered = listRegisteredProjectStores();
expect(registered).toEqual(expect.arrayContaining([
{ projectId: "proj_list_a", store: storeA },
{ projectId: "proj_list_b", store: storeB },
]));
});
it("notifies on project store registration", async () => {
const listener = vi.fn();
const unsubscribe = onProjectStoreRegistered(listener);
const store = await getOrCreateProjectStore("proj_listener");
expect(listener).toHaveBeenCalledWith("proj_listener", store);
unsubscribe();
await getOrCreateProjectStore("proj_listener_2");
expect(listener).toHaveBeenCalledTimes(1);
});
it("reuses the same TaskStore instance for repeated calls with the same projectId", async () => {
const store1 = await getOrCreateProjectStore("proj_abc");
const store2 = await getOrCreateProjectStore("proj_abc");

View File

@@ -39,34 +39,61 @@ export function decideIssueAction(
}
export class GitHubTrackingStateService {
private readonly store: TaskStore;
private readonly onTaskMoved = (event: TaskMovedEvent): void => {
void this.handleTaskMoved(event);
};
private readonly onTaskDeleted = (task: Task, meta?: { githubIssueAction?: GithubIssueAction }): void => {
void this.handleTaskDeleted(task, meta);
};
private readonly defaultStore: TaskStore;
private readonly listeners = new Map<TaskStore, {
onTaskMoved: (event: TaskMovedEvent) => void;
onTaskDeleted: (task: Task, meta?: { githubIssueAction?: GithubIssueAction }) => void;
}>();
private started = false;
constructor(store: TaskStore) {
this.store = store;
this.defaultStore = store;
}
start(): void {
if (this.started) return;
this.started = true;
this.store.on("task:moved", this.onTaskMoved);
this.store.on("task:deleted", this.onTaskDeleted);
this.attach(this.defaultStore);
}
stop(): void {
if (!this.started) return;
this.started = false;
this.store.off("task:moved", this.onTaskMoved);
this.store.off("task:deleted", this.onTaskDeleted);
for (const store of this.listeners.keys()) {
this.detach(store);
}
}
private async handleTaskMoved(event: TaskMovedEvent): Promise<void> {
attach(store: TaskStore): void {
if (this.listeners.has(store)) {
return;
}
const onTaskMoved = (event: TaskMovedEvent): void => {
void this.handleTaskMoved(store, event);
};
const onTaskDeleted = (task: Task, meta?: { githubIssueAction?: GithubIssueAction }): void => {
void this.handleTaskDeleted(store, task, meta);
};
this.listeners.set(store, { onTaskMoved, onTaskDeleted });
if (this.started) {
store.on("task:moved", onTaskMoved);
store.on("task:deleted", onTaskDeleted);
}
}
detach(store: TaskStore): void {
const handlers = this.listeners.get(store);
if (!handlers) {
return;
}
store.off("task:moved", handlers.onTaskMoved);
store.off("task:deleted", handlers.onTaskDeleted);
this.listeners.delete(store);
}
private async handleTaskMoved(store: TaskStore, event: TaskMovedEvent): Promise<void> {
const decision = decideIssueAction(event.from, event.to);
if (!decision) {
return;
@@ -83,7 +110,7 @@ export class GitHubTrackingStateService {
const { owner, repo, number } = issue;
if (!owner || !repo || !number) {
await this.store.logEntry(
await store.logEntry(
event.task.id,
"Failed to update GitHub tracking issue state",
"Linked issue metadata is incomplete",
@@ -92,11 +119,11 @@ export class GitHubTrackingStateService {
}
try {
const projectSettings = await this.store.getSettings() as Pick<ProjectSettings, "githubAuthMode" | "githubAuthToken">;
const globalSettings = (await this.store.getGlobalSettingsStore?.()?.getSettings?.() ?? {}) as Pick<GlobalSettings, never>;
const projectSettings = await store.getSettings() as Pick<ProjectSettings, "githubAuthMode" | "githubAuthToken">;
const globalSettings = (await store.getGlobalSettingsStore?.()?.getSettings?.() ?? {}) as Pick<GlobalSettings, never>;
const resolution = resolveGithubTrackingAuth({ projectSettings, globalSettings });
if (!resolution.ok) {
await this.store.logEntry(event.task.id, "Skipped GitHub tracking issue state update", resolution.message);
await store.logEntry(event.task.id, "Skipped GitHub tracking issue state update", resolution.message);
return;
}
@@ -111,7 +138,7 @@ export class GitHubTrackingStateService {
decision.action === "close" ? "closed" : "open",
decision.stateReason,
);
await this.store.logEntry(
await store.logEntry(
event.task.id,
decision.action === "close"
? "Closed linked GitHub tracking issue"
@@ -119,7 +146,7 @@ export class GitHubTrackingStateService {
`${owner}/${repo}#${number}`,
);
} catch (err) {
await this.store.logEntry(
await store.logEntry(
event.task.id,
decision.action === "close"
? "Failed to close GitHub tracking issue"
@@ -129,7 +156,7 @@ export class GitHubTrackingStateService {
}
}
private async handleTaskDeleted(task: Task, meta?: { githubIssueAction?: GithubIssueAction }): Promise<void> {
private async handleTaskDeleted(store: TaskStore, task: Task, meta?: { githubIssueAction?: GithubIssueAction }): Promise<void> {
if (task.githubTracking?.enabled !== true) {
return;
}
@@ -146,12 +173,12 @@ export class GitHubTrackingStateService {
const githubIssueAction = meta?.githubIssueAction ?? "auto";
if (githubIssueAction === "leave") {
await this.store.logEntry(task.id, "Left linked GitHub tracking issue unchanged on task delete", `${owner}/${repo}#${number}`);
await store.logEntry(task.id, "Left linked GitHub tracking issue unchanged on task delete", `${owner}/${repo}#${number}`);
return;
}
const projectSettings = await this.store.getSettings() as Pick<ProjectSettings, "githubAuthMode" | "githubAuthToken">;
const globalSettings = (await this.store.getGlobalSettingsStore?.()?.getSettings?.() ?? {}) as Pick<GlobalSettings, never>;
const projectSettings = await store.getSettings() as Pick<ProjectSettings, "githubAuthMode" | "githubAuthToken">;
const globalSettings = (await store.getGlobalSettingsStore?.()?.getSettings?.() ?? {}) as Pick<GlobalSettings, never>;
const resolution = resolveGithubTrackingAuth({ projectSettings, globalSettings });
if (!resolution.ok) {
return;
@@ -164,9 +191,9 @@ export class GitHubTrackingStateService {
if (githubIssueAction === "delete") {
try {
await client.deleteIssue(owner, repo, number);
await this.store.logEntry(task.id, "Deleted linked GitHub tracking issue", `${owner}/${repo}#${number}`);
await store.logEntry(task.id, "Deleted linked GitHub tracking issue", `${owner}/${repo}#${number}`);
} catch (err) {
await this.store.logEntry(task.id, "Failed to delete linked GitHub tracking issue", err instanceof Error ? err.message : String(err));
await store.logEntry(task.id, "Failed to delete linked GitHub tracking issue", err instanceof Error ? err.message : String(err));
}
return;
}
@@ -174,7 +201,7 @@ export class GitHubTrackingStateService {
try {
await client.setIssueState(owner, repo, number, "closed", "not_planned");
} catch (err) {
await this.store.logEntry(task.id, "Failed to close linked GitHub tracking issue", err instanceof Error ? err.message : String(err));
await store.logEntry(task.id, "Failed to close linked GitHub tracking issue", err instanceof Error ? err.message : String(err));
}
}
}

View File

@@ -36,6 +36,7 @@ const pendingCreations = new Map<string, Promise<TaskStore>>();
* lookups.
*/
const initializedProjects = new Set<string>();
const projectRegisteredListeners = new Set<(projectId: string, store: TaskStore) => void>();
/**
* Optional callback invoked once when a new project store is first created.
@@ -102,6 +103,10 @@ export async function getOrCreateProjectStore(projectId: string): Promise<TaskSt
_onProjectFirstCreated(projectId);
}
for (const listener of projectRegisteredListeners) {
listener(projectId, store);
}
return store;
})();
@@ -148,3 +153,14 @@ export function invalidateAllGlobalSettingsCaches(): void {
store.getGlobalSettingsStore().invalidateCache();
}
}
export function listRegisteredProjectStores(): Array<{ projectId: string; store: TaskStore }> {
return Array.from(storeCache.entries(), ([projectId, store]) => ({ projectId, store }));
}
export function onProjectStoreRegistered(listener: (projectId: string, store: TaskStore) => void): () => void {
projectRegisteredListeners.add(listener);
return () => {
projectRegisteredListeners.delete(listener);
};
}

View File

@@ -24,6 +24,7 @@ import { GitHubIssueCommentService } from "../github-issue-comment.js";
import { GitHubTrackingCommentService } from "../github-tracking-comments.js";
import { GitHubTrackingStateService } from "../github-tracking-state.js";
import { githubRateLimiter } from "../github-poll.js";
import { listRegisteredProjectStores, onProjectStoreRegistered } from "../project-store-resolver.js";
import {
classifyWebhookEvent,
getGitHubAppConfig,
@@ -1126,7 +1127,32 @@ export function registerGitGitHubRoutes(ctx: ApiRoutesContext): void {
const githubTrackingStateService = new GitHubTrackingStateService(store);
githubTrackingStateService.start();
ctx.registerDispose(() => githubTrackingStateService.stop());
const attachedStateStores = new Set<TaskStore>();
const attachStateStore = (projectStore: TaskStore) => {
if (attachedStateStores.has(projectStore)) {
return;
}
attachedStateStores.add(projectStore);
githubTrackingStateService.attach(projectStore);
};
attachStateStore(store);
for (const { store: projectStore } of listRegisteredProjectStores()) {
attachStateStore(projectStore);
}
const unsubscribeProjectStoreRegistration = onProjectStoreRegistered((_projectId, projectStore) => {
attachStateStore(projectStore);
});
ctx.registerDispose(() => {
unsubscribeProjectStoreRegistration();
for (const projectStore of attachedStateStores) {
githubTrackingStateService.detach(projectStore);
}
githubTrackingStateService.stop();
});
}
/**