feat(web-next): bindAll resolves IEventBus + IJobQueue per env
resolveEventsAndJobsProduction wires PayloadJobsEventBus + PayloadJobQueue against the bootstrapped Payload instance; resolveEventsAndJobsDevSeed wires InMemoryEventBus + InMemoryJobQueue. Each per-feature binder now receives (bus, queue) so Phase 7 generators can subscribe handlers and register job tasks at the <gen:event-handlers>/<gen:jobs> anchors.
This commit is contained in:
@@ -17,6 +17,7 @@
|
|||||||
"@repo/blog": "workspace:*",
|
"@repo/blog": "workspace:*",
|
||||||
"@repo/core-api": "workspace:*",
|
"@repo/core-api": "workspace:*",
|
||||||
"@repo/core-cms": "workspace:*",
|
"@repo/core-cms": "workspace:*",
|
||||||
|
"@repo/core-events": "workspace:*",
|
||||||
"@repo/core-shared": "workspace:*",
|
"@repo/core-shared": "workspace:*",
|
||||||
"@repo/core-trpc": "workspace:*",
|
"@repo/core-trpc": "workspace:*",
|
||||||
"@repo/core-ui": "workspace:*",
|
"@repo/core-ui": "workspace:*",
|
||||||
|
|||||||
@@ -1,6 +1,7 @@
|
|||||||
import { describe, it, expect, vi, beforeEach, afterEach } from "vitest";
|
import { describe, it, expect, vi, beforeEach, afterEach } from "vitest";
|
||||||
|
|
||||||
vi.mock("@repo/core-cms", () => ({ default: Promise.resolve({}) }));
|
vi.mock("@repo/core-cms", () => ({ default: Promise.resolve({}) }));
|
||||||
|
vi.mock("payload", () => ({ getPayload: vi.fn(async () => ({ jobs: { queue: vi.fn() } })) }));
|
||||||
vi.mock("@repo/blog/di/bind-production", () => ({ bindProductionBlog: vi.fn() }));
|
vi.mock("@repo/blog/di/bind-production", () => ({ bindProductionBlog: vi.fn() }));
|
||||||
vi.mock("@repo/auth/di/bind-production", () => ({ bindProductionAuth: vi.fn() }));
|
vi.mock("@repo/auth/di/bind-production", () => ({ bindProductionAuth: vi.fn() }));
|
||||||
vi.mock("@repo/marketing-pages/di/bind-production", () => ({ bindProductionMarketingPages: vi.fn() }));
|
vi.mock("@repo/marketing-pages/di/bind-production", () => ({ bindProductionMarketingPages: vi.fn() }));
|
||||||
@@ -50,6 +51,41 @@ describe("bindAllProduction", () => {
|
|||||||
await bindAllProduction();
|
await bindAllProduction();
|
||||||
expect(bindProductionBlog).toHaveBeenCalledOnce();
|
expect(bindProductionBlog).toHaveBeenCalledOnce();
|
||||||
});
|
});
|
||||||
|
|
||||||
|
it("passes a Payload-backed bus + queue to each per-feature binder", async () => {
|
||||||
|
const { bindAllProduction } = await import("./bind-production");
|
||||||
|
const { bindProductionAuth } = await import("@repo/auth/di/bind-production");
|
||||||
|
const { PayloadJobsEventBus } = await import("@repo/core-events");
|
||||||
|
const { PayloadJobQueue } = await import("@repo/core-shared/jobs");
|
||||||
|
|
||||||
|
await bindAllProduction();
|
||||||
|
|
||||||
|
const args = vi.mocked(bindProductionAuth).mock.calls[0]!;
|
||||||
|
// Args: (config, tracer, logger, bus, queue)
|
||||||
|
expect(args[3]).toBeInstanceOf(PayloadJobsEventBus);
|
||||||
|
expect(args[4]).toBeInstanceOf(PayloadJobQueue);
|
||||||
|
});
|
||||||
|
});
|
||||||
|
|
||||||
|
describe("bindAllDevSeed", () => {
|
||||||
|
beforeEach(() => {
|
||||||
|
vi.resetModules();
|
||||||
|
vi.clearAllMocks();
|
||||||
|
});
|
||||||
|
|
||||||
|
it("passes an in-memory bus + queue to each per-feature dev-seed binder", async () => {
|
||||||
|
const { bindAllDevSeed } = await import("./bind-production");
|
||||||
|
const { bindDevSeedAuth } = await import("@repo/auth/di/bind-dev-seed");
|
||||||
|
const { InMemoryEventBus } = await import("@repo/core-events");
|
||||||
|
const { InMemoryJobQueue } = await import("@repo/core-shared/jobs");
|
||||||
|
|
||||||
|
await bindAllDevSeed();
|
||||||
|
|
||||||
|
const args = vi.mocked(bindDevSeedAuth).mock.calls[0]!;
|
||||||
|
// Args: (tracer, logger, bus, queue)
|
||||||
|
expect(args[2]).toBeInstanceOf(InMemoryEventBus);
|
||||||
|
expect(args[3]).toBeInstanceOf(InMemoryJobQueue);
|
||||||
|
});
|
||||||
});
|
});
|
||||||
|
|
||||||
describe("bindAll dispatcher", () => {
|
describe("bindAll dispatcher", () => {
|
||||||
|
|||||||
@@ -2,6 +2,7 @@
|
|||||||
// SERVER-ONLY: this module imports Payload config and must never be bundled into the browser.
|
// SERVER-ONLY: this module imports Payload config and must never be bundled into the browser.
|
||||||
import "reflect-metadata";
|
import "reflect-metadata";
|
||||||
import { Container } from "inversify";
|
import { Container } from "inversify";
|
||||||
|
import { getPayload } from "payload";
|
||||||
import config from "@repo/core-cms";
|
import config from "@repo/core-cms";
|
||||||
import {
|
import {
|
||||||
bindNoopInstrumentation,
|
bindNoopInstrumentation,
|
||||||
@@ -9,6 +10,16 @@ import {
|
|||||||
type ITracer,
|
type ITracer,
|
||||||
type ILogger,
|
type ILogger,
|
||||||
} from "@repo/core-shared/instrumentation";
|
} from "@repo/core-shared/instrumentation";
|
||||||
|
import {
|
||||||
|
InMemoryEventBus,
|
||||||
|
PayloadJobsEventBus,
|
||||||
|
type IEventBus,
|
||||||
|
} from "@repo/core-events";
|
||||||
|
import {
|
||||||
|
InMemoryJobQueue,
|
||||||
|
PayloadJobQueue,
|
||||||
|
type IJobQueue,
|
||||||
|
} from "@repo/core-shared/jobs";
|
||||||
import { bindProductionBlog } from "@repo/blog/di/bind-production";
|
import { bindProductionBlog } from "@repo/blog/di/bind-production";
|
||||||
import { bindProductionAuth } from "@repo/auth/di/bind-production";
|
import { bindProductionAuth } from "@repo/auth/di/bind-production";
|
||||||
import { bindProductionMarketingPages } from "@repo/marketing-pages/di/bind-production";
|
import { bindProductionMarketingPages } from "@repo/marketing-pages/di/bind-production";
|
||||||
@@ -29,6 +40,8 @@ const sharedContainer = new Container();
|
|||||||
|
|
||||||
let resolvedTracer: ITracer | null = null;
|
let resolvedTracer: ITracer | null = null;
|
||||||
let resolvedLogger: ILogger | null = null;
|
let resolvedLogger: ILogger | null = null;
|
||||||
|
let resolvedBus: IEventBus | null = null;
|
||||||
|
let resolvedQueue: IJobQueue | null = null;
|
||||||
|
|
||||||
/** Rule 0: pick instrumentation backend from DSN env (orthogonal to repo mode). */
|
/** Rule 0: pick instrumentation backend from DSN env (orthogonal to repo mode). */
|
||||||
function resolveInstrumentation(): { tracer: ITracer; logger: ILogger } {
|
function resolveInstrumentation(): { tracer: ITracer; logger: ILogger } {
|
||||||
@@ -44,6 +57,41 @@ function resolveInstrumentation(): { tracer: ITracer; logger: ILogger } {
|
|||||||
return result;
|
return result;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Production-mode event bus + job queue: backed by Payload's job system so
|
||||||
|
* events fan out via durable Payload tasks (`PayloadJobsEventBus` enqueues
|
||||||
|
* `__events.<publisher>.<event>.<consumer>` per subscribed handler) and ad-hoc
|
||||||
|
* jobs go through `PayloadJobQueue.enqueue`. Cached after first resolution.
|
||||||
|
*/
|
||||||
|
async function resolveEventsAndJobsProduction(): Promise<{
|
||||||
|
bus: IEventBus;
|
||||||
|
queue: IJobQueue;
|
||||||
|
}> {
|
||||||
|
if (resolvedBus && resolvedQueue) return { bus: resolvedBus, queue: resolvedQueue };
|
||||||
|
const resolvedConfig = await config;
|
||||||
|
const payload = await getPayload({ config: resolvedConfig });
|
||||||
|
const queue = new PayloadJobQueue(payload);
|
||||||
|
const bus = new PayloadJobsEventBus(queue);
|
||||||
|
resolvedBus = bus;
|
||||||
|
resolvedQueue = queue;
|
||||||
|
return { bus, queue };
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Dev-seed mode: in-process bus + queue. Per-feature binders register their
|
||||||
|
* job handlers via `queue.register(slug, handler)` and subscribe their event
|
||||||
|
* handlers via `bus.subscribe(...)` at bind time, so dev/test exercises the
|
||||||
|
* full publish → handler → enqueue path without booting Payload.
|
||||||
|
*/
|
||||||
|
function resolveEventsAndJobsDevSeed(): { bus: IEventBus; queue: IJobQueue } {
|
||||||
|
if (resolvedBus && resolvedQueue) return { bus: resolvedBus, queue: resolvedQueue };
|
||||||
|
const queue = new InMemoryJobQueue();
|
||||||
|
const bus = new InMemoryEventBus();
|
||||||
|
resolvedBus = bus;
|
||||||
|
resolvedQueue = queue;
|
||||||
|
return { bus, queue };
|
||||||
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Production path: swap each feature's mock repository binding for the real
|
* Production path: swap each feature's mock repository binding for the real
|
||||||
* Payload-backed one. Constructs `new XRepository(config, tracer, logger)` per
|
* Payload-backed one. Constructs `new XRepository(config, tracer, logger)` per
|
||||||
@@ -53,12 +101,13 @@ export async function bindAllProduction(): Promise<void> {
|
|||||||
if (bound) return;
|
if (bound) return;
|
||||||
bound = true;
|
bound = true;
|
||||||
const { tracer, logger } = resolveInstrumentation(); // Rule 0
|
const { tracer, logger } = resolveInstrumentation(); // Rule 0
|
||||||
|
const { bus, queue } = await resolveEventsAndJobsProduction();
|
||||||
const resolvedConfig = await config;
|
const resolvedConfig = await config;
|
||||||
bindProductionAuth(resolvedConfig, tracer, logger); // Phase E task 19
|
bindProductionAuth(resolvedConfig, tracer, logger, bus, queue); // Phase E task 19
|
||||||
bindProductionBlog(resolvedConfig, tracer, logger); // Phase E task 18
|
bindProductionBlog(resolvedConfig, tracer, logger, bus, queue); // Phase E task 18
|
||||||
bindProductionMarketingPages(resolvedConfig, tracer, logger); // Phase E task 20
|
bindProductionMarketingPages(resolvedConfig, tracer, logger, bus, queue); // Phase E task 20
|
||||||
bindProductionNavigation(resolvedConfig, tracer, logger); // Phase E task 21
|
bindProductionNavigation(resolvedConfig, tracer, logger, bus, queue); // Phase E task 21
|
||||||
bindProductionMedia(resolvedConfig, tracer, logger); // Phase E task 22
|
bindProductionMedia(resolvedConfig, tracer, logger, bus, queue); // Phase E task 22
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
@@ -70,11 +119,12 @@ export async function bindAllDevSeed(): Promise<void> {
|
|||||||
if (bound) return;
|
if (bound) return;
|
||||||
bound = true;
|
bound = true;
|
||||||
const { tracer, logger } = resolveInstrumentation(); // Rule 0
|
const { tracer, logger } = resolveInstrumentation(); // Rule 0
|
||||||
await bindDevSeedAuth(tracer, logger); // Phase E task 19
|
const { bus, queue } = resolveEventsAndJobsDevSeed();
|
||||||
await bindDevSeedBlog(tracer, logger); // Phase E task 18
|
await bindDevSeedAuth(tracer, logger, bus, queue); // Phase E task 19
|
||||||
await bindDevSeedMarketingPages(tracer, logger); // Phase E task 20
|
await bindDevSeedBlog(tracer, logger, bus, queue); // Phase E task 18
|
||||||
await bindDevSeedNavigation(tracer, logger); // Phase E task 21
|
await bindDevSeedMarketingPages(tracer, logger, bus, queue); // Phase E task 20
|
||||||
await bindDevSeedMedia(tracer, logger); // Phase E task 22
|
await bindDevSeedNavigation(tracer, logger, bus, queue); // Phase E task 21
|
||||||
|
await bindDevSeedMedia(tracer, logger, bus, queue); // Phase E task 22
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
@@ -106,6 +156,8 @@ export function __resetBindStateForTests(): void {
|
|||||||
bound = false;
|
bound = false;
|
||||||
resolvedTracer = null;
|
resolvedTracer = null;
|
||||||
resolvedLogger = null;
|
resolvedLogger = null;
|
||||||
|
resolvedBus = null;
|
||||||
|
resolvedQueue = null;
|
||||||
}
|
}
|
||||||
|
|
||||||
/** Test-only accessor for resolved instrumentation. */
|
/** Test-only accessor for resolved instrumentation. */
|
||||||
|
|||||||
3
pnpm-lock.yaml
generated
3
pnpm-lock.yaml
generated
@@ -154,6 +154,9 @@ importers:
|
|||||||
'@repo/core-cms':
|
'@repo/core-cms':
|
||||||
specifier: workspace:*
|
specifier: workspace:*
|
||||||
version: link:../../packages/core-cms
|
version: link:../../packages/core-cms
|
||||||
|
'@repo/core-events':
|
||||||
|
specifier: workspace:*
|
||||||
|
version: link:../../packages/core-events
|
||||||
'@repo/core-shared':
|
'@repo/core-shared':
|
||||||
specifier: workspace:*
|
specifier: workspace:*
|
||||||
version: link:../../packages/core-shared
|
version: link:../../packages/core-shared
|
||||||
|
|||||||
Reference in New Issue
Block a user