Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
35 commits
Select commit Hold shift + click to select a range
f65694a
wip(relay): hold webhook requests for environments that opt in
juliusmarminge Oct 4, 2026
0fafa46
feat(relay,server,web,mobile): opt-in to hold webhooks while offline
juliusmarminge Oct 4, 2026
0e9d168
fix(server): relay deliveries run once and rate-limited ones can retry
juliusmarminge Oct 4, 2026
7e1db6e
refactor(relay,server): hold webhooks in a Durable Object that pushes…
juliusmarminge Oct 4, 2026
c4897cc
fix: held webhooks keep their arrival time and invalid age limits are…
juliusmarminge Oct 4, 2026
a860f90
fix(relay,server): held webhooks keep their receive time and one cap …
juliusmarminge Oct 4, 2026
fca6d77
fix(mobile): a deleted webhook task no longer shows its URL or Rotate
juliusmarminge Oct 5, 2026
22f8bcb
fix(mobile): Rotate URL is disabled while rotating, saving, or discon…
juliusmarminge Oct 5, 2026
0756745
fix(web,mobile): an invalid webhook age limit blocks saving on both c…
juliusmarminge Oct 5, 2026
7e569d9
refactor(client-runtime): share the webhook prompt default and age-li…
juliusmarminge Oct 5, 2026
4227c75
fix(client-runtime): new webhook tasks start with a body-only prompt
juliusmarminge Oct 5, 2026
9784b3f
fix(web): settings search offers webhook holding only when its row shows
juliusmarminge Oct 5, 2026
51c9791
fix(web): thread automations hide Run now for webhook tasks
juliusmarminge Oct 5, 2026
8fc6767
docs(user): webhook URLs need a managed tunnel, and what mobile can't do
juliusmarminge Oct 5, 2026
7e98750
fix(server): relay webhook URLs use the tunnel key, not the environme…
juliusmarminge Oct 5, 2026
ad90b99
fix(mobile): the webhook trigger is labelled On webhook, as on web
juliusmarminge Oct 5, 2026
eeaabef
perf(web): skip the T3 Connect link-state read in builds without T3 C…
juliusmarminge Oct 5, 2026
577fa27
fix(relay): webhook URLs name one managed endpoint, not a claimable e…
juliusmarminge Oct 5, 2026
f4ed062
fix(server): a webhook request dropped mid-flight no longer loses its…
juliusmarminge Oct 5, 2026
5086c1f
fix(server): a relay receive time in the future counts as now
juliusmarminge Oct 5, 2026
9d1f984
fix(server): webhook prompts redact credential headers like the deliv…
juliusmarminge Oct 5, 2026
51cdc7b
fix(server): webhook prompts past the provider input limit are not di…
juliusmarminge Oct 5, 2026
f697367
fix(server): webhook route answers defects with the fixed 500 body
juliusmarminge Oct 5, 2026
1eeb161
fix(server): one tunnel reconnect wakes held webhooks once
juliusmarminge Oct 5, 2026
794b942
fix(relay): held webhooks retry every 10 s for 3 minutes after a wake
juliusmarminge Oct 5, 2026
659081c
feat(relay,server): webhooks and held deliveries can be monitored
juliusmarminge Oct 5, 2026
844ae83
fix(relay): held webhooks can't be jammed, inherited, or delayed by a…
juliusmarminge Oct 5, 2026
f120f71
docs(relay): inbox run details live on the store spans
juliusmarminge Oct 5, 2026
6eea59a
fix(server): a webhook prompt's length is counted as the provider cou…
juliusmarminge Oct 5, 2026
4898478
fix(relay): held requests behind busy hooks' backlogs are reached in …
juliusmarminge Oct 5, 2026
6d33885
feat(relay,server): the relay records what the environment did with e…
juliusmarminge Oct 5, 2026
00775e0
fix(relay): the worker's own request span never records a webhook token
juliusmarminge Oct 5, 2026
bbe3a6b
fix(relay): a failed delivery run cannot push back a wake's earlier a…
juliusmarminge Oct 5, 2026
614183b
feat(web): a webhook task says whether T3 Connect forwards or holds i…
juliusmarminge Oct 5, 2026
734a2fe
feat(relay,server): the environment trusts relay delivery headers onl…
juliusmarminge Oct 5, 2026
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
Original file line number Diff line number Diff line change
Expand Up @@ -4,14 +4,21 @@ import type {
ScheduledTask,
ScheduledTaskUpsertInput,
} from "@t3tools/contracts";
import { resolveEnvironmentMachineKind } from "@t3tools/contracts";
import {
MAX_WEBHOOK_DELIVERY_AGE_MINUTES,
resolveEnvironmentMachineKind,
} from "@t3tools/contracts";
import type { MenuAction } from "@react-native-menu/menu";
import { DateTimePicker } from "@expo/ui/community/datetime-picker";
import {
isAtomCommandInterrupted,
squashAtomCommandFailure,
type AtomCommandResult,
} from "@t3tools/client-runtime/state/runtime";
import {
DEFAULT_WEBHOOK_PROMPT,
parseMaxDeliveryAge,
} from "@t3tools/client-runtime/scheduled-task-webhook";
import {
useCallback,
useEffect,
Expand Down Expand Up @@ -59,7 +66,6 @@ import { SettingsSection } from "./components/SettingsSection";
import { useSettingsEnvironmentFilter, type SettingsTarget } from "./settings-environment-filter";
import {
editDraft,
DEFAULT_WEBHOOK_PROMPT,
scheduledTaskDefaultModel,
scheduleFromDraft,
type ScheduledTaskDraft as Draft,
Expand Down Expand Up @@ -132,7 +138,7 @@ function FormField(props: {
readonly label: string;
readonly value: string;
readonly onChange: (value: string) => void;
readonly keyboardType?: "decimal-pad";
readonly keyboardType?: "decimal-pad" | "number-pad";
readonly disabled?: boolean;
readonly placeholder?: string;
readonly borderTop?: boolean;
Expand Down Expand Up @@ -611,6 +617,16 @@ function TaskForm({
? { ...draft.schedule, signature: liveTask.schedule.signature }
: draft.schedule,
);
if (
draft.schedule.mode === "webhook" &&
parseMaxDeliveryAge(draft.schedule.maxDeliveryAgeMinutes) === undefined
) {
Alert.alert(
"Invalid age limit",
`Enter whole minutes from 1 to ${MAX_WEBHOOK_DELIVERY_AGE_MINUTES}, or leave it blank.`,
);
return;
}
if (
!draft.title.trim() ||
!draft.prompt.trim() ||
Expand Down Expand Up @@ -823,7 +839,7 @@ function TaskForm({
options={[
{ value: "fixed_time", label: "At a time" },
{ value: "interval", label: "Interval" },
{ value: "webhook", label: "Webhook" },
{ value: "webhook", label: "On webhook" },
]}
selected={draft.schedule.mode}
onSelect={(mode) => {
Expand Down Expand Up @@ -910,13 +926,31 @@ function TaskForm({
/>
</>
) : draft.schedule.mode === "webhook" ? (
<WebhookScheduleDetails
environmentId={environmentId}
task={
tasks.data?.tasks.find((task) => task.id === draft.task?.id) ?? draft.task ?? null
}
signatureConfigured={draft.schedule.signature !== null}
/>
<>
<WebhookScheduleDetails
environmentId={environmentId}
// The live row, so a rotated URL shows up without reopening the form.
// Once the list has loaded, a missing task is gone; don't keep showing its URL.
task={
tasks.data
? (tasks.data.tasks.find((task) => task.id === draft.task?.id) ?? null)
: draft.task
}
signatureConfigured={draft.schedule.signature !== null}
disabled={saving || environmentUnavailable}
/>
<FormField
label="Skip requests older than (minutes)"
value={draft.schedule.maxDeliveryAgeMinutes}
placeholder="Run every request"
keyboardType="number-pad"
disabled={saving}
borderTop
onChange={(maxDeliveryAgeMinutes) =>
setDraft({ ...draft, schedule: { ...draft.schedule, maxDeliveryAgeMinutes } })
}
/>
</>
) : (
<>
<FormField
Expand Down Expand Up @@ -976,11 +1010,14 @@ function WebhookScheduleDetails({
environmentId,
task,
signatureConfigured,
disabled,
}: {
readonly environmentId: EnvironmentId;
readonly task: ScheduledTask | null;
readonly signatureConfigured: boolean;
readonly disabled: boolean;
}) {
const [rotating, setRotating] = useState(false);
const rotate = useAtomCommand(serverEnvironment.rotateScheduledTaskWebhookToken, {
label: "scheduled task rotate webhook token",
reportFailure: false,
Expand Down Expand Up @@ -1023,27 +1060,34 @@ function WebhookScheduleDetails({
) : null}
<Pressable
accessibilityRole="button"
accessibilityState={{ disabled: disabled || rotating }}
disabled={disabled || rotating}
onPress={() =>
Alert.alert("Rotate URL?", "The current URL stops working immediately.", [
{ text: "Cancel", style: "cancel" },
{
text: "Rotate",
style: "destructive",
onPress: () =>
onPress: () => {
setRotating(true);
void rotate({ environmentId, input: { id: task.id } }).then((result) => {
setRotating(false);
if (result._tag === "Failure" && !isAtomCommandInterrupted(result)) {
Alert.alert(
"Could not rotate URL",
String(squashAtomCommandFailure(result)),
);
}
}),
});
},
},
])
}
className="min-h-11 justify-center active:opacity-70"
className="min-h-11 justify-center active:opacity-70 disabled:opacity-50"
>
<Text className="text-base text-danger-foreground">Rotate URL</Text>
<Text className="text-base text-danger-foreground">
{rotating ? "Rotating…" : "Rotate URL"}
</Text>
</Pressable>
</>
)}
Expand Down
22 changes: 20 additions & 2 deletions apps/mobile/src/features/settings/scheduledTaskDraft.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -32,7 +32,25 @@ describe("scheduleDraftForTask", () => {
it("round-trips a webhook schedule without a signature", () => {
const draft = scheduleDraftForTask({ schedule: { type: "webhook", signature: null } });
expect(draft.mode).toBe("webhook");
expect(scheduleFromDraft(draft)).toEqual({ type: "webhook", signature: null });
expect(scheduleFromDraft(draft)).toEqual({
type: "webhook",
signature: null,
maxDeliveryAgeMinutes: null,
});
});

it("round-trips a webhook max age and treats blank input as no limit", () => {
const draft = scheduleDraftForTask({
schedule: { type: "webhook", signature: null, maxDeliveryAgeMinutes: 45 },
});
expect(draft.maxDeliveryAgeMinutes).toBe("45");
expect(scheduleFromDraft(draft)).toMatchObject({ maxDeliveryAgeMinutes: 45 });
// An invalid limit is an invalid schedule, never a silently removed one.
expect(scheduleFromDraft({ ...draft, maxDeliveryAgeMinutes: "1.5" })).toBeNull();
expect(scheduleFromDraft({ ...draft, maxDeliveryAgeMinutes: "0" })).toBeNull();
expect(scheduleFromDraft({ ...draft, maxDeliveryAgeMinutes: "" })).toMatchObject({
maxDeliveryAgeMinutes: null,
});
});

it("keeps a webhook signature on save without sending a secret", () => {
Expand All @@ -44,7 +62,7 @@ describe("scheduleDraftForTask", () => {
const saved = scheduleFromDraft(
scheduleDraftForTask({ schedule: { type: "webhook", signature } }),
);
expect(saved).toEqual({ type: "webhook", signature });
expect(saved).toEqual({ type: "webhook", signature, maxDeliveryAgeMinutes: null });
expect(saved?.type === "webhook" && saved.signature && "secret" in saved.signature).toBe(false);
});
});
Expand Down
21 changes: 17 additions & 4 deletions apps/mobile/src/features/settings/scheduledTaskDraft.ts
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@ import type {
} from "@t3tools/contracts";

import { DEFAULT_SERVER_SETTINGS } from "@t3tools/contracts";
import { parseMaxDeliveryAge } from "@t3tools/client-runtime/scheduled-task-webhook";
import {
resolveProjectSettings,
type LegacyProjectSettingsFields,
Expand Down Expand Up @@ -44,6 +45,8 @@ export type ScheduleDraft = {
readonly intervalMinutes: string;
/** A webhook signature check configured elsewhere; mobile keeps it but does not edit it. */
readonly signature: ScheduledTaskWebhookSignature | null;
/** Minutes as typed; empty runs every held request regardless of age. */
readonly maxDeliveryAgeMinutes: string;
};

export const DEFAULT_SCHEDULE: ScheduleDraft = {
Expand All @@ -52,11 +55,9 @@ export const DEFAULT_SCHEDULE: ScheduleDraft = {
weekdays: [1, 2, 3, 4, 5],
intervalMinutes: "15",
signature: null,
maxDeliveryAgeMinutes: "",
};

/** Prompt a new webhook task starts with: the whole request, which the user can narrow down. */
export const DEFAULT_WEBHOOK_PROMPT = "Handle this webhook:\n{{request}}";

export function scheduleDraftForTask(task: Pick<ScheduledTask, "schedule">): ScheduleDraft {
switch (task.schedule.type) {
case "fixed_time":
Expand All @@ -74,12 +75,22 @@ export function scheduleDraftForTask(task: Pick<ScheduledTask, "schedule">): Sch
intervalMinutes: String(Math.max(1, task.schedule.everyMs / 60_000)),
};
case "webhook":
return { ...DEFAULT_SCHEDULE, mode: "webhook", signature: task.schedule.signature };
return {
...DEFAULT_SCHEDULE,
mode: "webhook",
signature: task.schedule.signature,
maxDeliveryAgeMinutes:
task.schedule.maxDeliveryAgeMinutes == null
? ""
: String(task.schedule.maxDeliveryAgeMinutes),
};
}
}

export function scheduleFromDraft(draft: ScheduleDraft): ScheduledTaskUpsertSchedule | null {
if (draft.mode === "webhook") {
const maxDeliveryAgeMinutes = parseMaxDeliveryAge(draft.maxDeliveryAgeMinutes);
if (maxDeliveryAgeMinutes === undefined) return null;
// No secret is sent, so the server keeps the stored one.
return {
type: "webhook",
Expand All @@ -91,6 +102,7 @@ export function scheduleFromDraft(draft: ScheduleDraft): ScheduledTaskUpsertSche
encoding: draft.signature.encoding,
prefix: draft.signature.prefix,
},
maxDeliveryAgeMinutes,
};
}
if (draft.mode === "interval") {
Expand Down Expand Up @@ -145,6 +157,7 @@ function draftSignature(draft: ScheduledTaskDraft): string {
draft.schedule.timeOfDay,
[...draft.schedule.weekdays].sort((a, b) => a - b),
draft.schedule.intervalMinutes,
draft.schedule.maxDeliveryAgeMinutes,
draft.workspace,
draft.baseRef,
draft.checkoutPath,
Expand Down
30 changes: 30 additions & 0 deletions apps/server/src/cloud/ManagedEndpointRuntime.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -209,6 +209,36 @@ describe("CloudManagedEndpointRuntime", () => {
}),
);

it.effect("signals each registered tunnel connection", () =>
Effect.gen(function* () {
const output = yield* Queue.unbounded<Uint8Array>();
const spawner = ChildProcessSpawner.make(() =>
Effect.gen(function* () {
const handle = makeHandle({
pid: 700,
onKill: () => {},
output: Stream.fromQueue(output),
});
yield* Effect.addFinalizer(() => handle.kill().pipe(Effect.ignore));
return handle;
}),
);
const runtime = yield* buildCloudManagedEndpointRuntime(spawner);
yield* runtime.applyConfig({
providerKind: "cloudflare_tunnel",
connectorToken: "token",
tunnelId: "tunnel-1",
});
yield* Queue.offer(
output,
new TextEncoder().encode(
"2026-10-04T06:30:43Z INF Registered tunnel connection connIndex=0 event=0\n",
),
);
expect(Option.isSome(yield* Stream.runHead(runtime.tunnelConnected))).toBe(true);
}),
);

it.effect("recovers a rejected tunnel without waiting for the connector to exit", () =>
Effect.gen(function* () {
const output = yield* Queue.unbounded<Uint8Array>();
Expand Down
9 changes: 8 additions & 1 deletion apps/server/src/cloud/ManagedEndpointRuntime.ts
Original file line number Diff line number Diff line change
Expand Up @@ -47,6 +47,8 @@ export class CloudManagedEndpointRuntime extends Context.Service<
) => Effect.Effect<CloudManagedEndpointRuntimeStatus>;
readonly recoveryRequests: Stream.Stream<RelayManagedEndpointRuntimeConfig>;
readonly requestRecovery: (config: RelayManagedEndpointRuntimeConfig) => Effect.Effect<void>;
/** Emits when the connector registers a tunnel connection, i.e. the relay can reach us again. */
readonly tunnelConnected: Stream.Stream<void>;
readonly withLinkStateLock: <A, E, R>(effect: Effect.Effect<A, E, R>) => Effect.Effect<A, E, R>;
}
>()("t3/cloud/ManagedEndpointRuntime/CloudManagedEndpointRuntime") {}
Expand Down Expand Up @@ -134,6 +136,7 @@ export const make = Effect.gen(function* () {
const activeRef = yield* Ref.make<ActiveConnector | null>(null);
const desiredConfigRef = yield* Ref.make<RelayManagedEndpointRuntimeConfig | null>(null);
const recoveryRequests = yield* Queue.sliding<RelayManagedEndpointRuntimeConfig>(1);
const tunnelConnections = yield* Queue.sliding<void>(1);
const reconcileSemaphore = yield* Semaphore.make(1);
const restartDelayRef = yield* Ref.make(0);
const linkStateSemaphore = yield* Semaphore.make(1);
Expand Down Expand Up @@ -236,7 +239,10 @@ export const make = Effect.gen(function* () {
switch (classifyRelayClientOutput(line)) {
case "connected":
rejectedRegistrations = 0;
return Effect.logInfo("Relay client tunnel connection registered", attributes);
return Effect.logInfo("Relay client tunnel connection registered", attributes).pipe(
Effect.andThen(Queue.offer(tunnelConnections, undefined)),
Effect.asVoid,
);
case "warning":
if (isRejectedRelayClientTunnelOutput(line)) {
rejectedRegistrations += 1;
Expand Down Expand Up @@ -412,6 +418,7 @@ export const make = Effect.gen(function* () {
applyConfig,
recoveryRequests: Stream.fromQueue(recoveryRequests),
requestRecovery: (config) => Queue.offer(recoveryRequests, config).pipe(Effect.asVoid),
tunnelConnected: Stream.fromQueue(tunnelConnections),
withLinkStateLock: linkStateSemaphore.withPermits(1),
});

Expand Down
31 changes: 31 additions & 0 deletions apps/server/src/cloud/config.ts
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@ export const RELAY_URL_SECRET = "cloud-relay-url";
export const RELAY_ISSUER_SECRET = "cloud-relay-issuer";
export const RELAY_ENVIRONMENT_CREDENTIAL_SECRET = "cloud-relay-environment-credential";
export const PUBLISH_AGENT_ACTIVITY_SECRET = "cloud-publish-agent-activity";
export const HOLD_WEBHOOKS_WHILE_OFFLINE_SECRET = "cloud-hold-webhooks-while-offline";

export const encodeEndpointRuntimeConfigJson = Schema.encodeEffect(
Schema.fromJsonString(RelayManagedEndpointRuntimeConfig),
Expand Down Expand Up @@ -73,3 +74,33 @@ export const readAgentActivityPublishingActive = (
environmentCredential !== ""
);
}).pipe(Effect.orElseSucceed(() => false));

const readSecretString = (
secrets: ServerSecretStore.ServerSecretStore["Service"],
name: string,
): Effect.Effect<string | null> =>
secrets.get(name).pipe(
Effect.map((bytes) =>
Option.isSome(bytes) && bytes.value.length > 0 ? new TextDecoder().decode(bytes.value) : null,
),
Effect.orElseSucceed(() => null),
Comment thread
coderabbitai[bot] marked this conversation as resolved.
);

/** The relay URL and environment credential, or null when not linked to T3 Connect. */
export const readRelayConnection = (secrets: ServerSecretStore.ServerSecretStore["Service"]) =>
Effect.all([
readSecretString(secrets, RELAY_URL_SECRET),
readSecretString(secrets, RELAY_ENVIRONMENT_CREDENTIAL_SECRET),
]).pipe(
Effect.map(([url, environmentCredential]) =>
url && environmentCredential ? { url, environmentCredential } : null,
),
);

/** Whether this environment opted in to T3 Connect holding webhooks while it is offline. */
export const readHoldWebhooksWhileOffline = (
secrets: ServerSecretStore.ServerSecretStore["Service"],
) =>
readSecretString(secrets, HOLD_WEBHOOKS_WHILE_OFFLINE_SECRET).pipe(
Effect.map((value) => value === "true"),
);
2 changes: 2 additions & 0 deletions apps/server/src/cloud/http.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -249,6 +249,7 @@ describe("reconcileDesiredCloudLink", () => {
applyConfig: unusedSecretStoreOperation,
recoveryRequests: Stream.empty,
requestRecovery: () => Effect.void,
tunnelConnected: Stream.empty,
withLinkStateLock: (effect) => effect,
} satisfies ManagedEndpointRuntime.CloudManagedEndpointRuntime["Service"]),
),
Expand Down Expand Up @@ -410,6 +411,7 @@ describe("releaseManagedTunnelOnShutdown", () => {
}),
recoveryRequests: Stream.empty,
requestRecovery: () => Effect.void,
tunnelConnected: Stream.empty,
withLinkStateLock: (effect) => effect,
}),
),
Expand Down
Loading
Loading