Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion apps/hub/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -21,7 +21,7 @@
"@corbits/mailbox": "github:corbitsdev/corbits-mailbox#65590a85fa143251b3ac20ba2eca92fc35e70e51",
"@corbits/memory": "github:corbitsdev/corbits-memory#e74da20f148a302dff5400915fe504ee2395e913",
"@corbits/url-path": "workspace:*",
"@corbits/webhooks": "github:corbitsdev/webhooks#570bd52688920c0ce018a8fbcd046199fa4125dd",
"@corbits/webhooks": "github:corbitsdev/webhooks#5c9e7d8fad13ebfbc747d01981ed4b0b44b810cd",
"@corbits/workflows": "workspace:*",
"@corbits/xai-provider": "github:corbitsdev/corbits-xai-provider#9f2d4bac40ea075df092fef807404a13c638cf14",
"@intx/authz": "0.3.0",
Expand Down
91 changes: 50 additions & 41 deletions apps/hub/src/server.ts
Original file line number Diff line number Diff line change
Expand Up @@ -51,21 +51,20 @@ import { Hono } from "hono";
// Everything above this line is upstream Interchange's server.ts, verbatim;
// see AGENTS.md's "plain Interchange tenant" ruling.
import {
buildMailFrame,
createInMemoryMailboxEventBus,
createMailboxDb,
createMailboxPersist,
generateMailboxMessageId,
mountMailbox,
} from "@corbits/mailbox";
import { createMemory, loadMemoryConfig } from "@corbits/memory";
import { createCronTicker, mountCron } from "@corbits/cron";
import { createCronTicker, createRunTriggerCronDeliver, mountCron } from "@corbits/cron";
import {
createHubMailboxAuthorizeSender,
createHubPersistMailWithSessionEnsure,
} from "./mailbox-persist";
import { captureMailboxRequest, createMailboxDeliver } from "./mailbox-send";
import { installWebhooks, type HookMailRouter } from "@corbits/webhooks";
import { reportError } from "@corbits/error-sink";
import { createRunTriggerDeliverer, installWebhooks, type HookMailRouter } from "@corbits/webhooks";
import {
createProcessSidecarProvisioner,
readProcessProvisionerConfig,
Expand Down Expand Up @@ -546,6 +545,22 @@ export async function createHubServer({
app.route(`${TENANT_PREFIX}/mailbox`, mailboxApp);
}

// The router every system-originated trigger (webhook, cron) goes
// through. `HookMailRouter` types its payloads as `unknown` at the
// package boundary; this just narrows them back to `sidecarRouter`'s own
// types on the way through, with no behavior change.
const systemTriggerMailRouter: HookMailRouter = {
routeMail: (address, rawMessage, authenticatedSender, messageId) =>
sidecarRouter.routeMail(address, rawMessage, authenticatedSender, messageId),
sendRunGrants: (address, runId, stepGrants, senderIdentities) =>
sidecarRouter.sendRunGrants(
address,
runId,
stepGrants as Parameters<typeof sidecarRouter.sendRunGrants>[2],
senderIdentities as Parameters<typeof sidecarRouter.sendRunGrants>[3],
),
};

let cronTicker: { start(): void; stop(): void } | undefined;
{
const cronApp = new Hono<TenantEnv>();
Expand All @@ -561,28 +576,37 @@ export async function createHubServer({
cronTicker = createCronTicker({
db,
intervalMs: 60_000,
senderAddressFor: (tenantId) => `cron@${tenantId}`,
deliver: async (message) => {
const tenantId = message.from.slice("cron@".length);
const [tenantRow] = await db
.select({ domain: tenantTable.domain })
.from(tenantTable)
.where(eq(tenantTable.id, tenantId))
.limit(1);
if (tenantRow === undefined) {
throw new Error(`no tenant "${tenantId}" to address cron mail from`);
}
const from = `cron@${tenantRow.domain}`;
await mailboxLookups.persistMail({
senderAddress: from,
recipients: message.to,
raw: buildMailFrame({
from,
to: message.to.join(", "),
subject: message.subject,
body: message.body,
messageId: generateMailboxMessageId(from),
// A due schedule fires with nobody signed in, so it cannot ride the
// mailbox persist path, which authorizes its sender against a live
// routable endpoint. It is a system trigger like an inbound webhook,
// so it takes the same route: the run is the authenticated sender of
// its own signed trigger mail, with its grants materialized first.
deliver: createRunTriggerCronDeliver(
createRunTriggerDeliverer({
router: systemTriggerMailRouter,
materialize: createMailTriggeredRunGrantsMaterializer({
db,
principalKeyStore,
grantStore,
}),
tenantDomain: async (tenantId) => {
const [tenantRow] = await db
.select({ domain: tenantTable.domain })
.from(tenantTable)
.where(eq(tenantTable.id, tenantId))
.limit(1);
if (tenantRow === undefined) {
throw new Error(`no tenant "${tenantId}" to address cron mail from`);
}
return tenantRow.domain;
},
senderLocalPart: "cron",
}),
),
onDeliveryError: (error, schedule) => {
reportError(error, {
operation: "hub.cron.deliver",
extra: { scheduleId: schedule.id, tenantId: schedule.tenantId },
});
},
});
Expand All @@ -600,27 +624,12 @@ export async function createHubServer({
app.route("/", memoryApp);
}

// `HookMailRouter` types its payloads as `unknown` at the package
// boundary; this just narrows them back to `sidecarRouter`'s own types
// on the way through, with no behavior change.
const webhookMailRouter: HookMailRouter = {
routeMail: (address, rawMessage, authenticatedSender, messageId) =>
sidecarRouter.routeMail(address, rawMessage, authenticatedSender, messageId),
sendRunGrants: (address, runId, stepGrants, senderIdentities) =>
sidecarRouter.sendRunGrants(
address,
runId,
stepGrants as Parameters<typeof sidecarRouter.sendRunGrants>[2],
senderIdentities as Parameters<typeof sidecarRouter.sendRunGrants>[3],
),
};

await installWebhooks({
app,
db,
credentialCipher,
principalKeyStore,
router: webhookMailRouter,
router: systemTriggerMailRouter,
});

// End of Corbits mount block.
Expand Down
4 changes: 2 additions & 2 deletions bun.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

31 changes: 31 additions & 0 deletions packages/cron/src/deliver.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,31 @@
import { expect, test } from "bun:test";

import { createRunTriggerCronDeliver } from "./deliver";

test("every recipient of a due schedule gets its own trigger", async () => {
const calls: Array<[string, string, string, string | undefined]> = [];
const deliver = createRunTriggerCronDeliver({
to: async (address, content, tenantId, subject) => {
calls.push([address, content, tenantId, subject]);
},
});

await deliver({ to: ["run-a@x", "run-b@x"], subject: "Daily", body: "go", tenantId: "t1" });

expect(calls).toEqual([
["run-a@x", "go", "t1", "Daily"],
["run-b@x", "go", "t1", "Daily"],
]);
});

test("a failed recipient stops the fan-out so the ticker can report it", async () => {
const deliver = createRunTriggerCronDeliver({
to: async (address) => {
if (address === "run-a@x") throw new Error("not routable");
},
});

expect(
deliver({ to: ["run-a@x", "run-b@x"], subject: "s", body: "b", tenantId: "t1" }),
).rejects.toThrow("not routable");
});
21 changes: 21 additions & 0 deletions packages/cron/src/deliver.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,21 @@
import type { DeliverCronMail } from "./ticker";

/** A system-trigger deliverer, shaped like `@corbits/webhooks`'s
* `MailDeliverer`. Structural so this package stays free of it. */
export type RunTriggerDeliverer = {
to: (
address: string,
content: string,
tenantId: string,
subject: string | undefined,
) => Promise<void>;
};

/** Fan a due schedule's recipients out over a run-trigger deliverer. */
export function createRunTriggerCronDeliver(deliverer: RunTriggerDeliverer): DeliverCronMail {
return async (message) => {
for (const address of message.to) {
await deliverer.to(address, message.body, message.tenantId, message.subject);
}
};
}
10 changes: 2 additions & 8 deletions packages/cron/src/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -10,12 +10,6 @@ export {
type ZonedParts,
} from "./cron";
export { cronScheduleTable, applyCronMigrations } from "./schema";
export {
createCronTicker,
cronSenderAddress,
type CronDb,
type CronSenderAddress,
type CronTicker,
type DeliverCronMail,
} from "./ticker";
export { createCronTicker, type CronDb, type CronTicker, type DeliverCronMail } from "./ticker";
export { mountCron, type MountCronOpts, type RequireTenantMember } from "./mount";
export { createRunTriggerCronDeliver, type RunTriggerDeliverer } from "./deliver";
39 changes: 23 additions & 16 deletions packages/cron/src/ticker.ts
Original file line number Diff line number Diff line change
Expand Up @@ -11,25 +11,25 @@ import { cronScheduleTable } from "./schema";
export type CronDb<TSchema extends Record<string, unknown> = Record<string, unknown>> =
PostgresJsDatabase<TSchema>;

/** A due schedule handed to the host. The sender identity is the host's to
* decide: only the host knows which addresses its mail transport
* authorizes, so this package names the tenant and never invents an
* address for it. */
export type DeliverCronMail = (message: {
to: string[];
subject: string;
body: string;
from: string;
tenantId: string;
}) => Promise<void> | void;

export type CronSenderAddress = (tenantId: string) => string;

/** The system sender address a tenant's cron mail comes from. */
export const cronSenderAddress: CronSenderAddress = (tenantId) => `cron@${tenantId}.internal`;

export type CreateCronTickerOpts<
TSchema extends Record<string, unknown> = Record<string, unknown>,
> = {
db: CronDb<TSchema>;
deliver: DeliverCronMail;
intervalMs: number;
senderAddressFor?: CronSenderAddress;
/** Told about a delivery that failed, so the host can report it. */
onDeliveryError?: (error: unknown, schedule: { id: string; tenantId: string }) => void;
};

export type CronTicker = {
Expand All @@ -56,7 +56,7 @@ function isDue(
async function tick<TSchema extends Record<string, unknown>>(
db: CronDb<TSchema>,
deliver: DeliverCronMail,
senderAddressFor: CronSenderAddress,
onDeliveryError: (error: unknown, schedule: { id: string; tenantId: string }) => void,
) {
await db.transaction(async (tx) => {
const now = new Date();
Expand All @@ -70,12 +70,19 @@ async function tick<TSchema extends Record<string, unknown>>(
.for("update", { skipLocked: true });

for (const row of candidates.filter((row) => isDue(row, now))) {
await deliver({
to: [row.toAddress],
subject: row.subject,
body: row.body,
from: senderAddressFor(row.tenantId),
});
// One schedule's failed delivery is its own: the tick still advances
// every due row, so a permanently undeliverable schedule cannot block
// the rest of the table or re-fire every minute forever.
try {
await deliver({
to: [row.toAddress],
subject: row.subject,
body: row.body,
tenantId: row.tenantId,
});
} catch (error) {
onDeliveryError(error, { id: row.id, tenantId: row.tenantId });
}
await tx
.update(cronScheduleTable)
.set({ lastFiredAt: now })
Expand All @@ -88,13 +95,13 @@ async function tick<TSchema extends Record<string, unknown>>(
export function createCronTicker<TSchema extends Record<string, unknown>>(
opts: CreateCronTickerOpts<TSchema>,
): CronTicker {
const senderAddressFor = opts.senderAddressFor ?? cronSenderAddress;
const onDeliveryError = opts.onDeliveryError ?? (() => undefined);
let timer: ReturnType<typeof setInterval> | undefined;
let inFlight: Promise<void> | undefined;

const runTick = () => {
if (inFlight !== undefined) return;
inFlight = tick(opts.db, opts.deliver, senderAddressFor).finally(() => {
inFlight = tick(opts.db, opts.deliver, onDeliveryError).finally(() => {
inFlight = undefined;
});
};
Expand Down
Loading