mirror of
https://github.com/VibedByKaKi/t3-code-android-nightly.git
synced 2026-10-09 11:51:15 +02:00
152 lines
6.6 KiB
TypeScript
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),
|
|
),
|
|
};
|
|
});
|
|
}),
|
|
);
|