t3-code-android-nightly/infra/relay/src/hooks/HookInboxObject.ts
Julius Marminge 8bd7745f0a
refactor(relay): import HookInboxObject once, as a namespace (#16307)
Co-authored-by: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
2026-10-05 19:30:02 -07:00

152 lines
6.6 KiB
TypeScript

import * as SqliteClient from "@effect/sql-sqlite-do/SqliteClient";
import type * as Alchemy from "alchemy";
import * as Cloudflare from "alchemy/Cloudflare";
import * as Clock from "effect/Clock";
import * as Effect from "effect/Effect";
import * as Option from "effect/Option";
import * as Result from "effect/Result";
import * as FetchHttpClient from "effect/http/FetchHttpClient";
import * as HookInboxStore from "./HookInboxStore.ts";
import { sendUpstream, TUNNEL_OFFLINE_STATUS } from "./upstream.ts";
/** Statuses that mean the environment is not there; the whole inbox waits and backs off. */
const UNREACHABLE_STATUSES = new Set([502, 503, 504, TUNNEL_OFFLINE_STATUS]);
/**
* Statuses that mean the environment is there but did not take this request
* yet: its task's queue is full (429) or it failed while handling it (500).
* The environment drops a delivery id it has already run, so retrying is safe.
*/
const BUSY_STATUSES = new Set([429, 500]);
/** When a run itself fails, the next one is tried after this long. */
const RUN_FAILURE_RETRY_MS = 60_000;
type Call<A> = Effect.Effect<A, never, Alchemy.RuntimeContext>;
/**
* One per managed endpoint, addressed by endpoint key. Holds webhook requests
* the environment could not take, in SQLite, and pushes them back through
* its tunnel from the alarm, oldest first, backing off while it stays away.
*/
export class HookInboxObject extends Cloudflare.DurableObject<
HookInboxObject,
{
/** Holds a request for `baseUrl`; false when the inbox is full and nothing was stored. */
readonly hold: (hook: HookInboxStore.HeldHook, baseUrl: string) => Call<boolean>;
/** The environment is back at `baseUrl`: deliver what is waiting now. */
readonly wake: (baseUrl: string) => Call<boolean>;
readonly clear: () => Call<void>;
}
>()("HookInboxObject") {}
/**
* Each call into an inbox is its own trace: the alarm has no parent, and a
* hold arrives over Durable Object RPC without the forwarding span. The
* object's name is the endpoint key, which the forward spans carry too.
*/
const withInboxSpan =
(name: string, inboxId: string) =>
<A, E, R>(effect: Effect.Effect<A, E, R>) =>
effect.pipe(
Effect.withSpan(name, { root: true, attributes: { "relay.hook.endpoint_key": inboxId } }),
);
const deliver = (baseUrl: string, hook: HookInboxStore.HeldHook) =>
sendUpstream(baseUrl, hook).pipe(
Effect.result,
Effect.map((result) => {
// Unreachable or timed out: a timeout may still have run it, and the
// environment drops a delivery id it has already seen.
if (Result.isFailure(result)) {
return { outcome: "unreachable" as const, reason: result.failure._tag };
}
if (Option.isNone(result.success))
return { outcome: "unreachable" as const, reason: "timeout" };
const status = result.success.value.status;
const upstreamOutcome = result.success.value.outcome;
const outcome: HookInboxStore.DeliveryOutcome = UNREACHABLE_STATUSES.has(status)
? "unreachable"
: BUSY_STATUSES.has(status)
? "busy"
: "delivered";
return { outcome, reason: `status ${status}`, upstreamOutcome };
}),
Effect.tap(({ upstreamOutcome }) =>
upstreamOutcome === undefined
? Effect.void
: Effect.annotateCurrentSpan({ "relay.hook.upstream_outcome": upstreamOutcome }),
),
Effect.tap(({ outcome, reason }) =>
outcome === "delivered"
? Effect.void
: Effect.logInfo("Held webhook not delivered yet", {
outcome,
reason,
deliveryId: hook.id,
}),
),
Effect.map(({ outcome }) => outcome),
Effect.provide(FetchHttpClient.layer),
);
export const layer = HookInboxObject.make(
Effect.gen(function* () {
const state = yield* Cloudflare.DurableObjectState;
// The init phase returns the per-instance Effect, which alchemy runs once
// per object; only that inner Effect may touch storage.
// @effect-diagnostics-next-line returnEffectInGen:off
return Effect.gen(function* () {
const layerSql = SqliteClient.layer({ storage: state.raw.storage });
// Inboxes are opened by endpoint key, so spans line up with the
// forward spans that held their requests.
const inboxId = state.raw.id.name ?? state.raw.id.toString();
const run = <A, E, R>(effect: Effect.Effect<A, E, R>) =>
effect.pipe(Effect.provide(layerSql), Effect.orDie);
yield* run(HookInboxStore.migrate);
/** Moves the alarm to `at`, unless one is already due sooner. */
const scheduleBy = (at: number) =>
Effect.gen(function* () {
const current = yield* state.storage.getAlarm();
if (current === null || current > at) yield* state.storage.setAlarm(at);
});
return {
hold: (hook: HookInboxStore.HeldHook, baseUrl: string) =>
Effect.gen(function* () {
const dueAt = yield* run(HookInboxStore.hold(hook, baseUrl));
yield* Effect.annotateCurrentSpan({ "relay.inbox.stored": dueAt !== null });
if (dueAt === null) return false;
yield* scheduleBy(dueAt);
return true;
}).pipe(withInboxSpan("relay.inbox.hold", inboxId)),
wake: (baseUrl: string) =>
Effect.gen(function* () {
const pending = yield* run(HookInboxStore.wake(baseUrl));
yield* Effect.annotateCurrentSpan({ "relay.inbox.pending": pending });
if (pending) yield* state.storage.setAlarm(yield* Clock.currentTimeMillis);
return pending;
}).pipe(withInboxSpan("relay.inbox.wake", inboxId)),
clear: () =>
Effect.gen(function* () {
yield* run(HookInboxStore.clear);
yield* state.storage.deleteAlarm();
}).pipe(withInboxSpan("relay.inbox.clear", inboxId)),
alarm: () =>
run(HookInboxStore.deliverDue(deliver)).pipe(
// A wake during this run may already have asked for an earlier
// run; a backoff must not push it out.
Effect.flatMap((nextAt) => (nextAt === null ? Effect.void : scheduleBy(nextAt))),
Effect.catchCause((cause) =>
Effect.logWarning("Held webhook delivery run failed", { cause }).pipe(
Effect.andThen(Effect.annotateCurrentSpan({ "relay.inbox.run_result": "failed" })),
Effect.andThen(Clock.currentTimeMillis),
Effect.flatMap((now) => scheduleBy(now + RUN_FAILURE_RETRY_MS)),
),
),
withInboxSpan("relay.inbox.deliver", inboxId),
),
};
});
}),
);