From 2aa968a357f035e6be805723f9dd04455ba7a88a Mon Sep 17 00:00:00 2001 From: Anton Niklasson Date: Sat, 11 Jul 2026 22:20:44 +0200 Subject: [PATCH] sync: expose application service API --- package.json | 2 +- packages/sync/README.md | 31 ++++ packages/sync/src/cli.ts | 31 ++-- packages/sync/src/index.test.ts | 65 ++----- packages/sync/src/index.ts | 68 ++------ packages/sync/src/models.ts | 48 +++++ .../sync/src/providers/github/normalize.ts | 39 +---- packages/sync/src/service.test.ts | 161 +++++++++++++++++ packages/sync/src/service.ts | 164 ++++++++++++++++++ 9 files changed, 457 insertions(+), 152 deletions(-) create mode 100644 packages/sync/README.md create mode 100644 packages/sync/src/models.ts create mode 100644 packages/sync/src/service.test.ts create mode 100644 packages/sync/src/service.ts diff --git a/package.json b/package.json index b6e5e79..5301d75 100644 --- a/package.json +++ b/package.json @@ -7,7 +7,7 @@ "dev:web": "mkdir -p .logs && concurrently --kill-others -n server,web -c blue,green \"pnpm --filter server dev 2>&1 | tee .logs/server.log\" \"pnpm --filter web dev 2>&1 | tee .logs/web.log\"", "demo": "DEMO=1 pnpm dev", "demo:web": "DEMO=1 pnpm dev:web", - "build": "pnpm --filter server build && pnpm --filter web build", + "build": "pnpm --filter sync build && pnpm --filter server build && pnpm --filter web build", "build:desktop": "pnpm build && pnpm --filter desktop build && pnpm --filter desktop package", "lint": "oxlint", "fmt": "oxfmt .", diff --git a/packages/sync/README.md b/packages/sync/README.md new file mode 100644 index 0000000..1c81d17 --- /dev/null +++ b/packages/sync/README.md @@ -0,0 +1,31 @@ +# Sync service + +`sync` exposes a small application API for consumers such as the dashboard +server. Cache files, SQLite rows, repositories, GitHub providers, and engine +construction are private implementation details. + +```ts +import { createSyncService, type SyncService } from "sync"; + +const sync: SyncService = createSyncService({ + onBackgroundError: console.error, +}); + +await sync.sync({ instanceId: "github-com", kind: "prs" }); +const authored = sync.listAuthoredPullRequests("github-com"); + +sync.start({ intervalMs: 25_000 }); +await sync.close(); +``` + +## Contract + +- Read methods return normalized domain objects, never database rows. +- `sync` is awaited and reports its cycle result to the caller. +- `requestSync` is fire-and-forget; failures go to `onBackgroundError` and are + contained by the service. +- `start` and `stop` control the recurring loop. +- `close` waits for the loop and requested syncs, then releases storage. It is + terminal and idempotent. +- Consumers depend on `SyncService`, so tests can provide an in-memory fake + without opening SQLite, reading config, starting timers, or calling GitHub. diff --git a/packages/sync/src/cli.ts b/packages/sync/src/cli.ts index a95593d..f5d206f 100644 --- a/packages/sync/src/cli.ts +++ b/packages/sync/src/cli.ts @@ -4,6 +4,7 @@ import { openCache, wipeCacheFile } from "./cache/open.js"; import { CACHE_SCHEMA_VERSION } from "./cache/schema.js"; import { type Repository, createSqliteRepository } from "./cache/store.js"; import { type SyncKind, createSyncEngine, printSummary } from "./engine.js"; +import { createSyncServiceFromDependencies } from "./service.js"; const USAGE = `ghd-sync — github-dashboard sync engine @@ -69,16 +70,16 @@ async function onceCommand(args: string[]): Promise { return 1; } - const { engine, close } = openEngine(); + const service = openService(); try { - const summary = await engine.runOnce({ - instance: values.instance, + const summary = await service.sync({ + instanceId: values.instance, kind, }); printSummary(summary); return 0; } finally { - close(); + await service.close(); } } @@ -105,8 +106,8 @@ async function loopCommand(args: string[]): Promise { const countdown = createCountdown(intervalMs); - const { engine, close } = openEngine(); - engine.start({ + const service = openService(); + service.start({ intervalMs, onCycle: (summary) => { countdown.stop(); @@ -124,10 +125,10 @@ async function loopCommand(args: string[]): Promise { try { await waitForSigint(); countdown.stop(); - await engine.stop(); + await service.stop(); return 0; } finally { - close(); + await service.close(); } } @@ -175,7 +176,7 @@ function createCountdown(intervalMs: number): { }; } -function openEngine() { +function openService() { const { db, path, wiped } = openCache(); if (wiped) { process.stderr.write( @@ -183,8 +184,16 @@ function openEngine() { ); } const repo = createSqliteRepository(db); - const engine = createSyncEngine({ repo }); - return { db, engine, close: () => db.close() }; + return createSyncServiceFromDependencies({ + repo, + engine: createSyncEngine({ repo }), + closeStorage: () => db.close(), + onBackgroundError: (error) => { + process.stderr.write( + `background sync failed: ${error instanceof Error ? error.message : String(error)}\n`, + ); + }, + }); } function statusCommand(): number { diff --git a/packages/sync/src/index.test.ts b/packages/sync/src/index.test.ts index 4e92091..c94ea14 100644 --- a/packages/sync/src/index.test.ts +++ b/packages/sync/src/index.test.ts @@ -2,68 +2,33 @@ import { mkdtempSync, rmSync } from "node:fs"; import { tmpdir } from "node:os"; import { join } from "node:path"; import { afterEach, beforeEach, describe, expect, test } from "vitest"; -// Import only from the public surface — anything missing here means callers -// (notably the server in Phase B of #72) would have to reach into internals. -import { - CACHE_SCHEMA_VERSION, - type GitHubInstance, - type Repository, - type SyncEngine, - createSqliteRepository, - createSyncEngine, - openCache, - reconcileInstances, -} from "./index.js"; +import * as publicApi from "./index.js"; +import { createSyncService, type SyncService } from "./index.js"; describe("public surface", () => { let cacheRoot: string; - let prevXdg: string | undefined; + let previousXdg: string | undefined; + let service: SyncService | undefined; beforeEach(() => { cacheRoot = mkdtempSync(join(tmpdir(), "ghd-public-")); - prevXdg = process.env.XDG_CACHE_HOME; + previousXdg = process.env.XDG_CACHE_HOME; process.env.XDG_CACHE_HOME = cacheRoot; }); - afterEach(() => { - if (prevXdg === undefined) delete process.env.XDG_CACHE_HOME; - else process.env.XDG_CACHE_HOME = prevXdg; + afterEach(async () => { + await service?.close(); + if (previousXdg === undefined) delete process.env.XDG_CACHE_HOME; + else process.env.XDG_CACHE_HOME = previousXdg; rmSync(cacheRoot, { recursive: true, force: true }); }); - test("composes from public exports the way the server would", () => { - const { db, path } = openCache(); - try { - expect(path).toContain("github-dashboard/cache.sqlite"); + test("exposes the service boundary without storage or engine constructors", async () => { + expect(Object.keys(publicApi)).toEqual(["createSyncService"]); - const repo: Repository = createSqliteRepository(db); - expect(repo.getSchemaVersion()).toBe(CACHE_SCHEMA_VERSION); - - const fakeInstances: GitHubInstance[] = [ - { - id: "github-com", - label: "github.com", - baseUrl: "https://api.github.com", - token: "redacted", - username: "u", - }, - ]; - const { added, removed } = reconcileInstances(repo, fakeInstances); - expect(added).toEqual(["github-com"]); - expect(removed).toEqual([]); - - expect(repo.listInstanceIds()).toEqual(["github-com"]); - expect(repo.getPrPayloads("github-com", "authored")).toEqual([]); - expect(repo.listNotifications("github-com")).toEqual([]); - - const engine: SyncEngine = createSyncEngine({ - repo, - // Inject fake config so we don't hit the real ~/.config or network. - loadInstances: async () => fakeInstances, - }); - expect(engine.isRunning()).toBe(false); - } finally { - db.close(); - } + service = createSyncService(); + expect(service.listAuthoredPullRequests("github-com")).toEqual([]); + expect(service.listReviewRequests("github-com")).toEqual([]); + expect(service.listNotifications("github-com")).toEqual([]); }); }); diff --git a/packages/sync/src/index.ts b/packages/sync/src/index.ts index 468712d..5d420a9 100644 --- a/packages/sync/src/index.ts +++ b/packages/sync/src/index.ts @@ -1,56 +1,16 @@ -// Public surface of the sync package. Consumers (the server, the CLI, future -// adopters) should import only from here so the internal layout can change -// without breaking callers. -// -// Typical composition: -// const { db } = openCache(); -// const repo = createSqliteRepository(db); -// const engine = createSyncEngine({ repo }); -// await engine.runOnce(); -// engine.start({ intervalMs: 25_000, onCycle: console.log }); - -// Cache file (raw SQLite open with version check) + path utilities +// Public application boundary for embedding the sync engine. Storage, +// repositories, providers, and the engine itself are implementation details. +export type { Notification, PullRequest, ReviewSummary } from "./models.js"; export { - type Cache, - type OpenCacheResult, - openCache, - wipeCacheFile, -} from "./cache/open.js"; -export { CACHE_SCHEMA_VERSION } from "./cache/schema.js"; -export { resolveCachePath } from "./cache/path.js"; - -// Repository (storage contract + sqlite implementation) -export { - type InstanceRow, - type InstanceSummary, - type NotificationRow, - type PrKind, - type PrKindCount, - type PrRow, - type Repository, - type SyncStateRow, - createSqliteRepository, -} from "./cache/store.js"; - -// Config (reading ~/.config/github-dashboard/config.yml) -export { - type GitHubInstance, - instanceIdFromDomain, - loadInstances, - resolveConfigPath, -} from "./config.js"; - -// Sync engine -export { - type FetchSummary, - type InstanceResult, - type SyncCycleOptions, - type SyncCycleSummary, - type SyncEngine, - type SyncEngineDeps, - type SyncKind, - type SyncLoopOptions, - createSyncEngine, - printSummary, - reconcileInstances, + type CreateSyncServiceOptions, + type StartSyncOptions, + type SyncRequest, + type SyncService, + createSyncService, +} from "./service.js"; +export type { + FetchSummary, + InstanceResult, + SyncCycleSummary, + SyncKind, } from "./engine.js"; diff --git a/packages/sync/src/models.ts b/packages/sync/src/models.ts new file mode 100644 index 0000000..2bea9d5 --- /dev/null +++ b/packages/sync/src/models.ts @@ -0,0 +1,48 @@ +export type CiStatus = "success" | "failure" | "pending" | "unknown"; + +export interface ReviewSummary { + approved: string[]; + changesRequested: string[]; +} + +export interface PullRequest { + id: number | string; + number: number; + title: string; + body: string; + url: string; + repo: string; + createdAt: string; + updatedAt: string; + author: string; + authorAvatar: string; + draft: boolean; + ciStatus: CiStatus; + inMergeQueue: boolean; + autoMerge: boolean; + autoMergeAllowed: boolean; + headBranch: string; + baseBranch: string; + reviews: ReviewSummary; + reviewDecision: string | null; + mergeStateStatus: string | null; + unresolvedThreadCount: number; + additions: number; + deletions: number; + commits: number; + commentCount: number; + labels: string[]; + mergeable: boolean | null; + autoAssigned?: boolean; +} + +export interface Notification { + id: string; + title: string; + type: string | null; + reason: string; + repo: string; + url: string; + unread: boolean; + updatedAt: string; +} diff --git a/packages/sync/src/providers/github/normalize.ts b/packages/sync/src/providers/github/normalize.ts index 13fb92e..c310b12 100644 --- a/packages/sync/src/providers/github/normalize.ts +++ b/packages/sync/src/providers/github/normalize.ts @@ -1,6 +1,7 @@ import type { PrNode } from "./queries.js"; +import type { CiStatus, PullRequest, ReviewSummary } from "../../models.js"; -export type CiStatus = "success" | "failure" | "pending" | "unknown"; +export type { CiStatus, ReviewSummary } from "../../models.js"; export function mapCiStatus(state: string | null | undefined): CiStatus { switch (state) { @@ -23,11 +24,6 @@ export function mapMergeable(v: PrNode["mergeable"]): boolean | null { return null; } -export interface ReviewSummary { - approved: string[]; - changesRequested: string[]; -} - export function summarizeReviews( reviews: PrNode["reviews"]["nodes"], ): ReviewSummary { @@ -46,36 +42,7 @@ export function summarizeReviews( return { approved, changesRequested }; } -export interface NormalizedPr { - id: number | string; - number: number; - title: string; - body: string; - url: string; - repo: string; - createdAt: string; - updatedAt: string; - author: string; - authorAvatar: string; - draft: boolean; - ciStatus: CiStatus; - inMergeQueue: boolean; - autoMerge: boolean; - autoMergeAllowed: boolean; - headBranch: string; - baseBranch: string; - reviews: ReviewSummary; - reviewDecision: PrNode["reviewDecision"]; - mergeStateStatus: PrNode["mergeStateStatus"]; - unresolvedThreadCount: number; - additions: number; - deletions: number; - commits: number; - commentCount: number; - labels: string[]; - mergeable: boolean | null; - autoAssigned?: boolean; -} +export type NormalizedPr = PullRequest; export function normalizePr(node: PrNode): NormalizedPr { const ci = mapCiStatus( diff --git a/packages/sync/src/service.test.ts b/packages/sync/src/service.test.ts new file mode 100644 index 0000000..8c33b5c --- /dev/null +++ b/packages/sync/src/service.test.ts @@ -0,0 +1,161 @@ +import { mkdtempSync, rmSync } from "node:fs"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import { + afterEach, + beforeEach, + describe, + expect, + type Mock, + test, + vi, +} from "vitest"; +import { openCache } from "./cache/open.js"; +import { createSqliteRepository, type Repository } from "./cache/store.js"; +import type { SyncCycleSummary, SyncEngine } from "./engine.js"; +import { createSyncServiceFromDependencies } from "./service.js"; + +const emptySummary: SyncCycleSummary = { + startedAt: "2026-01-01T00:00:00.000Z", + finishedAt: "2026-01-01T00:00:00.000Z", + durationMs: 0, + results: [], +}; + +describe("SyncService", () => { + let cacheRoot: string; + let previousXdg: string | undefined; + let repo: Repository; + let closeDb: () => void; + let engine: SyncEngine; + let closeStorage: Mock<() => void>; + let onBackgroundError: Mock<(error: unknown) => void>; + + beforeEach(() => { + cacheRoot = mkdtempSync(join(tmpdir(), "ghd-service-")); + previousXdg = process.env.XDG_CACHE_HOME; + process.env.XDG_CACHE_HOME = cacheRoot; + + const { db } = openCache(); + repo = createSqliteRepository(db); + closeDb = () => db.close(); + engine = { + runOnce: vi.fn(async () => emptySummary), + start: vi.fn(), + stop: vi.fn(async () => {}), + isRunning: vi.fn(() => false), + }; + closeStorage = vi.fn(); + onBackgroundError = vi.fn(); + }); + + afterEach(() => { + closeDb(); + if (previousXdg === undefined) delete process.env.XDG_CACHE_HOME; + else process.env.XDG_CACHE_HOME = previousXdg; + rmSync(cacheRoot, { recursive: true, force: true }); + }); + + function createService() { + return createSyncServiceFromDependencies({ + repo, + engine, + closeStorage, + onBackgroundError, + }); + } + + test("returns normalized read models instead of storage rows", () => { + repo.upsertInstance({ + id: "github", + label: "GitHub", + baseUrl: "https://api.github.com", + username: "alice", + }); + repo.replaceNotifications("github", [ + { + instance_id: "github", + id: "7", + title: "Mention", + type: "Issue", + reason: "mention", + repo: "acme/widgets", + url: "https://github.com/acme/widgets/issues/1", + unread: 1, + updated_at: "2026-01-01T00:00:00.000Z", + }, + ]); + + expect(createService().listNotifications("github")).toEqual([ + { + id: "7", + title: "Mention", + type: "Issue", + reason: "mention", + repo: "acme/widgets", + url: "https://github.com/acme/widgets/issues/1", + unread: true, + updatedAt: "2026-01-01T00:00:00.000Z", + }, + ]); + }); + + test("delegates explicit sync and loop lifecycle to the engine", async () => { + const service = createService(); + const request = { instanceId: "github", kind: "prs" as const }; + + await expect(service.sync(request)).resolves.toBe(emptySummary); + expect(engine.runOnce).toHaveBeenCalledWith({ + instance: "github", + kind: "prs", + }); + + service.start({ intervalMs: 123 }); + expect(engine.start).toHaveBeenCalledWith({ intervalMs: 123 }); + + await service.stop(); + expect(engine.stop).toHaveBeenCalledOnce(); + }); + + test("contains failures from fire-and-forget sync requests", async () => { + const error = new Error("config unavailable"); + vi.mocked(engine.runOnce).mockRejectedValueOnce(error); + const service = createService(); + + service.requestSync({ instanceId: "github" }); + await service.stop(); + + expect(onBackgroundError).toHaveBeenCalledWith(error); + }); + + test("close waits for work, closes storage once, and rejects further use", async () => { + let finish!: () => void; + vi.mocked(engine.runOnce).mockReturnValueOnce( + new Promise((resolve) => { + finish = () => resolve(emptySummary); + }), + ); + const service = createService(); + service.requestSync(); + + const closing = service.close(); + expect(closeStorage).not.toHaveBeenCalled(); + finish(); + await closing; + + expect(closeStorage).toHaveBeenCalledOnce(); + await service.close(); + expect(closeStorage).toHaveBeenCalledOnce(); + expect(() => service.listNotifications("github")).toThrow( + "SyncService is closed", + ); + }); + + test("close releases storage even when stopping the engine fails", async () => { + vi.mocked(engine.stop).mockRejectedValueOnce(new Error("stop failed")); + const service = createService(); + + await expect(service.close()).rejects.toThrow("stop failed"); + expect(closeStorage).toHaveBeenCalledOnce(); + }); +}); diff --git a/packages/sync/src/service.ts b/packages/sync/src/service.ts new file mode 100644 index 0000000..19a4237 --- /dev/null +++ b/packages/sync/src/service.ts @@ -0,0 +1,164 @@ +import { openCache } from "./cache/open.js"; +import { type Repository, createSqliteRepository } from "./cache/store.js"; +import { + type SyncCycleSummary, + type SyncEngine, + type SyncKind, + createSyncEngine, +} from "./engine.js"; +import type { Notification, PullRequest } from "./models.js"; + +export interface SyncService { + listAuthoredPullRequests(instanceId: string): PullRequest[]; + listReviewRequests(instanceId: string): PullRequest[]; + listNotifications(instanceId: string): Notification[]; + + sync(request?: SyncRequest): Promise; + requestSync(request?: SyncRequest): void; + + start(options?: StartSyncOptions): void; + stop(): Promise; + close(): Promise; + isRunning(): boolean; +} + +export interface SyncRequest { + instanceId?: string; + kind?: SyncKind; +} + +export interface StartSyncOptions { + intervalMs?: number; + onCycle?: (summary: SyncCycleSummary) => void; + onError?: (error: unknown) => void; +} + +export interface CreateSyncServiceOptions { + /** Receives failures from fire-and-forget requestSync calls. */ + onBackgroundError?: (error: unknown) => void; +} + +interface SyncServiceDependencies { + repo: Repository; + engine: SyncEngine; + closeStorage: () => void; + onBackgroundError: (error: unknown) => void; +} + +export function createSyncService( + options: CreateSyncServiceOptions = {}, +): SyncService { + const { db } = openCache(); + const repo = createSqliteRepository(db); + const engine = createSyncEngine({ repo }); + + return createSyncServiceFromDependencies({ + repo, + engine, + // The public service intentionally hides cache diagnostics. Internal + // composition roots such as the CLI can use the dependency seam below. + closeStorage: () => db.close(), + onBackgroundError: + options.onBackgroundError ?? + ((error) => console.error("Background sync failed:", error)), + }); +} + +/** Internal composition seam used by focused service tests. */ +export function createSyncServiceFromDependencies( + dependencies: SyncServiceDependencies, +): SyncService { + const { repo, engine, closeStorage, onBackgroundError } = dependencies; + const pending = new Set>(); + let closed = false; + + function assertOpen(): void { + if (closed) throw new Error("SyncService is closed"); + } + + function listPullRequests( + instanceId: string, + kind: "authored" | "review_requested", + ): PullRequest[] { + assertOpen(); + return repo.getPrPayloads(instanceId, kind) as PullRequest[]; + } + + const engineOptions = (request: SyncRequest | undefined) => + request ? { instance: request.instanceId, kind: request.kind } : undefined; + + async function waitForPending(): Promise { + while (pending.size > 0) await Promise.all(pending); + } + + async function stop(): Promise { + if (closed) return; + await engine.stop(); + await waitForPending(); + } + + return { + listAuthoredPullRequests: (instanceId) => + listPullRequests(instanceId, "authored"), + listReviewRequests: (instanceId) => + listPullRequests(instanceId, "review_requested"), + listNotifications: (instanceId) => { + assertOpen(); + return repo.listNotifications(instanceId).map((row) => ({ + id: row.id, + title: row.title, + type: row.type, + reason: row.reason, + repo: row.repo, + url: row.url, + unread: row.unread === 1, + updatedAt: row.updated_at, + })); + }, + sync: (request) => { + assertOpen(); + return engine.runOnce(engineOptions(request)); + }, + requestSync: (request) => { + assertOpen(); + let tracked: Promise; + tracked = engine + .runOnce(engineOptions(request)) + .then( + () => undefined, + (error) => { + try { + onBackgroundError(error); + } catch (handlerError) { + console.error( + "Background sync error handler failed:", + handlerError, + ); + } + }, + ) + .finally(() => pending.delete(tracked)); + pending.add(tracked); + }, + start: (options) => { + assertOpen(); + engine.start(options); + }, + stop, + close: async () => { + if (closed) return; + closed = true; + const results = await Promise.allSettled([ + engine.stop(), + waitForPending(), + ]); + closeStorage(); + const failure = results.find( + (result): result is PromiseRejectedResult => + result.status === "rejected", + ); + if (failure) throw failure.reason; + }, + isRunning: () => !closed && engine.isRunning(), + }; +}