Skip to content

Commit dc4260a

Browse files
authored
Merge pull request #327 from PetalNet/refactor/console-effect-ser-idiom
refactor(console): extract terminal domain effect + call domain effects in-process (drop internal fetch, getRequestEvent, runRemote)
2 parents 0d1f81b + 7e10ffd commit dc4260a

62 files changed

Lines changed: 5091 additions & 4446 deletions

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.
Lines changed: 9 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,8 @@
1-
import { getRequestEvent } from "$app/server";
21
import { publicConfig } from "$lib/config";
3-
import { searchMockPalette, type PaletteSearchResponse } from "$lib/data/palette";
2+
import { searchMockPalette } from "$lib/data/palette";
3+
import { searchPalette } from "$lib/server/domain/palette/service";
4+
import { currentPrincipal } from "$lib/server/domain/principal";
5+
import { ConsoleDomain } from "$lib/server/domain/service";
46
import { Effect, Schema } from "effect";
57
import { Query } from "svelte-effect-runtime";
68

@@ -13,21 +15,11 @@ const input = Schema.Struct({
1315
* request; SvelteKit owns transport, validation, serialization, and cancellation semantics.
1416
*/
1517
export const searchCommandPalette = Query(input, ({ query: text }) =>
16-
Effect.promise(async () => {
18+
Effect.gen(function* () {
1719
if (publicConfig.dataMode === "mock") return searchMockPalette(text);
18-
19-
const event = getRequestEvent();
20-
const headers = new Headers({ accept: "application/json", origin: event.url.origin });
21-
for (const name of ["authorization", "cookie"] as const) {
22-
const value = event.request.headers.get(name);
23-
if (value) headers.set(name, value);
24-
}
25-
const base = publicConfig.consoleApiBase ?? `${event.url.origin}/api/v1`;
26-
const response = await event.fetch(
27-
`${base}/palette/search?q=${encodeURIComponent(text)}&limit=24`,
28-
{ headers },
29-
);
30-
if (!response.ok) throw new Error(`Palette search returned ${String(response.status)}`);
31-
return (await response.json()) as PaletteSearchResponse;
20+
const domain = yield* ConsoleDomain;
21+
const services = yield* domain.services;
22+
const principal = yield* currentPrincipal;
23+
return yield* searchPalette(services, principal, text, 24);
3224
}),
3325
);

apps/console/src/lib/components/CommandPalette.svelte

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,9 +1,9 @@
11
<script lang="ts">
22
import { goto } from "$app/navigation";
33
import { searchCommandPalette } from "$lib/command-palette.remote";
4-
import { runRemote } from "$lib/rpc/browser";
54
import type { PaletteItem, PaletteKind, PaletteSearchResponse } from "$lib/data/palette";
65
import { visibleNav } from "$lib/nav";
6+
import { Effect } from "effect";
77
import Icon from "./Icon.svelte";
88
import ModalSurface from "./ModalSurface.svelte";
99
@@ -107,7 +107,7 @@
107107
const current = ++requestId;
108108
timer = setTimeout(() => void (async () => {
109109
try {
110-
const result = await runRemote(searchCommandPalette({ query: queryText }));
110+
const result = await Effect.runPromise(searchCommandPalette({ query: queryText }));
111111
if (current === requestId) remote = result;
112112
} catch {
113113
if (current === requestId) failed = true;

apps/console/src/lib/data/cost.remote.ts

Lines changed: 6 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -23,10 +23,12 @@ export const compareCost = Query(costComparisonRequestSchema, (request) =>
2323
const domain = yield* ConsoleDomain;
2424
const services = yield* domain.services;
2525
const principal = yield* currentPrincipal;
26-
return yield* Effect.tryPromise({
27-
try: () => compareCostPair(services.db.app, principal.scopes, request, services.costMeter),
28-
catch: (cause) => cause,
29-
}).pipe(
26+
return yield* compareCostPair(
27+
services.db.app,
28+
principal.scopes,
29+
request,
30+
services.costMeter,
31+
).pipe(
3032
Effect.catch((cause) =>
3133
cause instanceof CostComparisonUnavailableError
3234
? HttpError("ServiceUnavailable", cause.message)

apps/console/src/lib/data/library.ts

Lines changed: 26 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -393,16 +393,22 @@ function mapLibraryItem(item: ApiLibraryItem, hold?: string): LibraryItemView {
393393
};
394394
}
395395

396+
type LiveLibraryResults = readonly [
397+
PromiseSettledResult<ApiEnvelope<ApiLibraryItem>>,
398+
PromiseSettledResult<ApiEnvelope<ApiLibraryLink>>,
399+
PromiseSettledResult<ApiEnvelope<ApiLibraryHold>>,
400+
PromiseSettledResult<ApiEnvelope<ApiLibraryCuration>>,
401+
PromiseSettledResult<ApiEnvelope<ApiLibraryCapability>>,
402+
];
403+
396404
/** Map the scope-filtered Rev3 read surface; optional sources fail independently and honestly. */
397-
export async function readLiveLibrary(fetchFn: typeof fetch = fetch): Promise<LibraryData> {
398-
const [itemsResult, linksResult, holdsResult, curationResult, capabilitiesResult] =
399-
await Promise.allSettled([
400-
readAllPages<ApiLibraryItem>("/library/items?limit=1000", fetchFn),
401-
readAllPages<ApiLibraryLink>("/library/links?limit=1000", fetchFn),
402-
readAllPages<ApiLibraryHold>("/library/holds?limit=1000", fetchFn),
403-
readAllPages<ApiLibraryCuration>("/library/curation?limit=1000", fetchFn),
404-
readAllPages<ApiLibraryCapability>("/library/capabilities?limit=1000", fetchFn),
405-
]);
405+
export function assembleLiveLibrary([
406+
itemsResult,
407+
linksResult,
408+
holdsResult,
409+
curationResult,
410+
capabilitiesResult,
411+
]: LiveLibraryResults): LibraryData {
406412
if (itemsResult.status === "rejected") return liveEmptyLibrary;
407413
const holds = holdsResult.status === "fulfilled" ? holdsResult.value.items : [];
408414
const holdByItem = new Map(
@@ -481,3 +487,14 @@ export async function readLiveLibrary(fetchFn: typeof fetch = fetch): Promise<Li
481487
},
482488
};
483489
}
490+
491+
export async function readLiveLibrary(fetchFn: typeof fetch = fetch): Promise<LibraryData> {
492+
const settled = await Promise.allSettled([
493+
readAllPages<ApiLibraryItem>("/library/items?limit=1000", fetchFn),
494+
readAllPages<ApiLibraryLink>("/library/links?limit=1000", fetchFn),
495+
readAllPages<ApiLibraryHold>("/library/holds?limit=1000", fetchFn),
496+
readAllPages<ApiLibraryCuration>("/library/curation?limit=1000", fetchFn),
497+
readAllPages<ApiLibraryCapability>("/library/capabilities?limit=1000", fetchFn),
498+
]);
499+
return assembleLiveLibrary(settled);
500+
}
Lines changed: 15 additions & 24 deletions
Original file line numberDiff line numberDiff line change
@@ -1,32 +1,23 @@
1-
import { getRequestEvent, query } from "$app/server";
2-
import type { WorkSettlementSnapshot } from "$lib/api/types";
31
import { publicConfig } from "$lib/config";
42
import { mockWorkSettlement } from "$lib/data/work-settlement";
5-
import { error } from "@sveltejs/kit";
3+
import { currentPrincipal } from "$lib/server/domain/principal";
4+
import { readWorkSettlement } from "$lib/server/domain/reads/work-settlement";
5+
import { ConsoleDomain } from "$lib/server/domain/service";
6+
import { Effect } from "effect";
7+
import { Query } from "svelte-effect-runtime";
68

79
function isMock(): boolean {
810
return publicConfig.dataMode === "mock";
911
}
1012

11-
function forwardedHeaders(): Headers {
12-
const incoming = getRequestEvent().request.headers;
13-
const headers = new Headers({ accept: "application/json" });
14-
for (const name of ["authorization", "cookie", "x-dev-principal"]) {
15-
const value = incoming.get(name);
16-
if (value) headers.set(name, value);
17-
}
18-
return headers;
19-
}
20-
2113
/** One caller-scoped RPC powers Work's settle strip and Library's task-history projection. */
22-
export const getWorkSettlement = query(async (): Promise<WorkSettlementSnapshot> => {
23-
if (isMock()) return mockWorkSettlement();
24-
const event = getRequestEvent();
25-
const base = publicConfig.consoleApiBase ?? `${event.url.origin}/api/v1`;
26-
const response = await event.fetch(`${base}/work/settlement`, {
27-
headers: forwardedHeaders(),
28-
});
29-
if (!response.ok)
30-
error(response.status, `Work settlement source returned ${String(response.status)}`);
31-
return (await response.json()) as WorkSettlementSnapshot;
32-
});
14+
export const getWorkSettlement = Query(
15+
Effect.gen(function* () {
16+
if (isMock()) return mockWorkSettlement();
17+
const domain = yield* ConsoleDomain;
18+
const services = yield* domain.services;
19+
const principal = yield* currentPrincipal;
20+
if (!services.tracker) return yield* Effect.die(new Error("Tracker is unavailable"));
21+
return yield* readWorkSettlement(services.tracker, principal.scopes);
22+
}),
23+
);

apps/console/src/lib/operations.remote.ts

Lines changed: 48 additions & 42 deletions
Original file line numberDiff line numberDiff line change
@@ -22,7 +22,7 @@ import type {
2222
TaskItem,
2323
WorkerItem,
2424
} from "$lib/api/types";
25-
import { executeOpPlane } from "$lib/server/api/console-api";
25+
import { executeOpPlane } from "$lib/server/domain/commands/op-plane";
2626
import { listDashboards } from "$lib/server/domain/dashboard/store";
2727
import { consumeOpRateLimit } from "$lib/server/domain/op-rate-limit";
2828
import { currentPrincipal } from "$lib/server/domain/principal";
@@ -37,6 +37,7 @@ import {
3737
readTasks as readTasksCore,
3838
} from "$lib/server/domain/reads/tracker-reads";
3939
import { ConsoleDomain } from "$lib/server/domain/service";
40+
import { readTerminalAccess as readTerminalAccessCore } from "$lib/server/domain/terminal/service";
4041
import { Effect } from "effect";
4142
import { Command, Query } from "svelte-effect-runtime";
4243

@@ -108,7 +109,8 @@ export const readPlaneRemote = Query("unchecked", (plane: ReadPlane) =>
108109
zookie: principal.zookie,
109110
} satisfies Me;
110111
if (plane === "health")
111-
return yield* Effect.tryPromise(async () => {
112+
// The bus-health probe is a single scoped lake read; a fault is a defect, not a caller error.
113+
return yield* Effect.promise(async () => {
112114
const rows = await services.db.admin<{ bus_heartbeat_at: string | null }[]>`
113115
select max(received_at)::text as bus_heartbeat_at
114116
from events where type = 'console.bus.health'`;
@@ -120,33 +122,31 @@ export const readPlaneRemote = Query("unchecked", (plane: ReadPlane) =>
120122
};
121123
});
122124
if (plane === "roster") {
123-
const result = yield* Effect.tryPromise(() =>
124-
readRosterCore(services.db.app, services.tracker, principal.scopes),
125-
);
126-
return { ...result, items: result.items.map((item) => flattenRosterItem(item as never)) };
125+
const result = yield* readRosterCore(services.db.app, services.tracker, principal.scopes);
126+
return { ...result, items: result.items.map((item) => flattenRosterItem(item)) };
127127
}
128-
if (plane === "executors")
129-
return yield* Effect.tryPromise(() => readExecutorsCore(services.db.app, principal.scopes));
128+
if (plane === "executors") return yield* readExecutorsCore(services.db.app, principal.scopes);
130129
if (plane === "tasks" || plane === "leases") {
131130
if (!services.tracker)
132131
return yield* Effect.die(new Error("Tracker read adapter is unavailable"));
133-
return plane === "tasks"
132+
return yield* plane === "tasks"
134133
? readTasksCore(services.tracker, principal.scopes)
135134
: readLeasesCore(services.tracker, principal.scopes);
136135
}
137136
if (plane === "dashboards")
138-
return yield* Effect.tryPromise(() =>
137+
// The canonical plane reads dashboards with a fixed limit and no cursor, so the only modeled
138+
// DashboardError (a bad cursor) is unreachable here — treat it as a defect, keeping this
139+
// shared read's channel empty for every consumer.
140+
return yield* Effect.orDie(
139141
listDashboards(services.db.app, principal.scopes, services.cursorSecret, { limit: 100 }),
140142
);
141143
if (plane === "catalog")
142-
return yield* Effect.tryPromise(() =>
143-
readEntity(services.db.app, principal.scopes, "registry", { limit: 1_000 }),
144-
);
144+
return yield* readEntity(services.db.app, principal.scopes, "registry", { limit: 1_000 });
145145
const kind = projectedKinds[plane];
146146
if (kind === "attention") {
147-
const envelope = yield* Effect.tryPromise(() =>
148-
readTypedEntity(services.db.app, principal.scopes, kind, { limit: 1_000 }),
149-
);
147+
const envelope = yield* readTypedEntity(services.db.app, principal.scopes, kind, {
148+
limit: 1_000,
149+
});
150150
// Attention items carry an operating lane; a caller only sees items for lanes it holds.
151151
return {
152152
...envelope,
@@ -158,11 +158,9 @@ export const readPlaneRemote = Query("unchecked", (plane: ReadPlane) =>
158158
};
159159
}
160160

161-
return yield* Effect.tryPromise(() =>
162-
kind === "subscription"
163-
? readTypedEntity(services.db.app, principal.scopes, kind, { limit: 1_000 })
164-
: readEntity(services.db.app, principal.scopes, kind, { limit: 1_000 }),
165-
);
161+
return yield* kind === "subscription"
162+
? readTypedEntity(services.db.app, principal.scopes, kind, { limit: 1_000 })
163+
: readEntity(services.db.app, principal.scopes, kind, { limit: 1_000 });
166164
}),
167165
);
168166

@@ -171,9 +169,16 @@ export const runStructuredQuery = Query("unchecked", (request: StructuredQuery)
171169
const domain = yield* ConsoleDomain;
172170
const services = yield* domain.services;
173171
const principal = yield* currentPrincipal;
174-
return yield* Effect.tryPromise(() =>
175-
runStructured(services.db.app, principal.scopes, request),
176-
);
172+
return yield* runStructured(services.db.app, principal.scopes, request);
173+
}),
174+
);
175+
176+
export const readTerminalAccessRemote = Query(
177+
Effect.gen(function* () {
178+
const domain = yield* ConsoleDomain;
179+
const services = yield* domain.services;
180+
const principal = yield* currentPrincipal;
181+
return yield* readTerminalAccessCore(services, principal);
177182
}),
178183
);
179184

@@ -203,25 +208,26 @@ export const executeNamedOp = Command(
203208
retryable: true,
204209
},
205210
} satisfies OpResult;
206-
// The one authoritative command plane: catalog lookup, arg validation, authz, proposal
207-
// posture, audit intent/outcome, and internal adapters — identical to POST /api/v1/op.
208-
const { body } = yield* Effect.tryPromise(() =>
209-
executeOpPlane(
210-
services,
211-
services.monitor,
212-
{
213-
schema_version: 1,
214-
id: input.id ?? crypto.randomUUID(),
215-
op: input.op,
216-
args: input.args,
217-
...(input.reason !== undefined ? { reason: input.reason } : {}),
218-
...(input.task_id !== undefined ? { task_id: input.task_id } : {}),
219-
dry_run: input.dry_run ?? false,
220-
},
221-
principal,
222-
),
211+
// The one authoritative command plane, composed in process from the domain: catalog lookup,
212+
// arg validation, authz, proposal posture, audit intent/outcome, and internal adapters —
213+
// identical to POST /api/v1/op, which now runs the SAME domain effect at its HTTP edge.
214+
const { body } = yield* executeOpPlane(
215+
services,
216+
services.monitor,
217+
{
218+
schema_version: 1,
219+
id: input.id ?? crypto.randomUUID(),
220+
op: input.op,
221+
args: input.args,
222+
...(input.reason !== undefined ? { reason: input.reason } : {}),
223+
...(input.task_id !== undefined ? { task_id: input.task_id } : {}),
224+
dry_run: input.dry_run ?? false,
225+
},
226+
principal,
223227
);
224-
return body as unknown as OpResult;
228+
// The command plane returns the loosely-typed HTTP op envelope; narrow it to its documented
229+
// OpResult shape (a genuine, unavoidable narrowing of the shared envelope, not a mock cast).
230+
return body as OpResult;
225231
}),
226232
);
227233

0 commit comments

Comments
 (0)