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
219 changes: 219 additions & 0 deletions apps/server/src/cloud/CloudPreferences.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,219 @@
import { EnvironmentId } from "@t3tools/contracts";
import { assert, it } from "@effect/vitest";
import * as Effect from "effect/Effect";
import * as Fiber from "effect/Fiber";
import * as Layer from "effect/Layer";
import * as Option from "effect/Option";
import * as FetchHttpClient from "effect/http/FetchHttpClient";

import * as ServerSecretStore from "../auth/ServerSecretStore.ts";
import * as ServerEnvironment from "../environment/ServerEnvironment.ts";
import * as AgentAwarenessRelay from "../relay/AgentAwarenessRelay.ts";
import * as CloudPreferences from "./CloudPreferences.ts";
import {
HOLD_WEBHOOKS_WHILE_OFFLINE_SECRET,
PUBLISH_AGENT_ACTIVITY_SECRET,
RELAY_ENVIRONMENT_CREDENTIAL_SECRET,
RELAY_URL_SECRET,
} from "./config.ts";

const encode = (value: string) => new TextEncoder().encode(value);

/** A linked environment whose secret store can refuse writes, and the relay calls it made. */
const withService = <A, E>(
options: {
readonly failHoldWrite?: boolean;
readonly failActivityWrite?: boolean;
readonly relayFails?: boolean;
readonly failActivityRead?: boolean;
readonly failHoldRead?: boolean;
/** The first relay call reports itself, then waits for this before answering. */
readonly holdFirstRelayCall?: {
readonly started: () => void;
readonly released: Promise<void>;
};
},
body: (input: {
readonly preferences: CloudPreferences.CloudPreferences["Service"];
readonly stored: Map<string, Uint8Array>;
readonly relayCalls: Array<boolean>;
}) => Effect.Effect<A, E, never>,
) =>
Effect.gen(function* () {
const stored = new Map<string, Uint8Array>([
[RELAY_URL_SECRET, encode("https://relay.test")],
[RELAY_ENVIRONMENT_CREDENTIAL_SECRET, encode("credential")],
[HOLD_WEBHOOKS_WHILE_OFFLINE_SECRET, encode("false")],
[PUBLISH_AGENT_ACTIVITY_SECRET, encode("false")],
]);
const relayCalls: Array<boolean> = [];
const fetch: typeof globalThis.fetch = Object.assign(
(_input: Parameters<typeof globalThis.fetch>[0], init?: RequestInit) => {
const payload = JSON.parse(String(init?.body)) as { holdWebhooksWhileOffline: boolean };
relayCalls.push(payload.holdWebhooksWhileOffline);
if (relayCalls.length === 1 && options.holdFirstRelayCall) {
options.holdFirstRelayCall.started();
return options.holdFirstRelayCall.released.then(() => Response.json(payload));
}
return options.relayFails
? Promise.resolve(new Response("unavailable", { status: 503 }))
: Promise.resolve(Response.json(payload));
},
{ preconnect: () => {} },
);
const dependencies = Layer.mergeAll(
Layer.mock(ServerSecretStore.ServerSecretStore)({
get: (name) =>
(options.failActivityRead && name === PUBLISH_AGENT_ACTIVITY_SECRET) ||
(options.failHoldRead && name === HOLD_WEBHOOKS_WHILE_OFFLINE_SECRET)
? Effect.fail(
new ServerSecretStore.SecretStoreReadError({
resource: name,
cause: new Error("busy"),
}),
)
: Effect.succeed(Option.fromNullishOr(stored.get(name))),
set: (name, value) =>
(options.failHoldWrite && name === HOLD_WEBHOOKS_WHILE_OFFLINE_SECRET) ||
(options.failActivityWrite && name === PUBLISH_AGENT_ACTIVITY_SECRET)
? Effect.fail(
new ServerSecretStore.SecretStorePersistError({
resource: name,
cause: new Error("read-only"),
}),
)
: Effect.sync(() => void stored.set(name, value)),
}),
Layer.mock(ServerEnvironment.ServerEnvironment)({
getEnvironmentId: Effect.succeed(EnvironmentId.make("environment-1")),
}),
Layer.mock(AgentAwarenessRelay.AgentAwarenessRelay)({ requestCatchUp: () => Effect.void }),
);
return yield* Effect.gen(function* () {
const preferences = yield* CloudPreferences.CloudPreferences;
return yield* body({ preferences, stored, relayCalls });
}).pipe(
Effect.provide(CloudPreferences.layer.pipe(Layer.provide(dependencies))),
Effect.provideService(FetchHttpClient.Fetch, fetch),
);
});

it.effect("tells the relay before saving the hold setting locally", () =>
withService({}, ({ preferences, stored, relayCalls }) =>
Effect.gen(function* () {
yield* preferences.update({ publishAgentActivity: true, holdWebhooksWhileOffline: true });
assert.deepEqual(relayCalls, [true]);
assert.equal(
new TextDecoder().decode(stored.get(HOLD_WEBHOOKS_WHILE_OFFLINE_SECRET)),
"true",
);
}),
),
);

it.effect("puts the relay back when the local save fails", () =>
withService({ failHoldWrite: true }, ({ preferences, stored, relayCalls }) =>
Effect.gen(function* () {
const error = yield* preferences
.update({ publishAgentActivity: true, holdWebhooksWhileOffline: true })
.pipe(Effect.flip);
assert.equal(error._tag, "EnvironmentHttpInternalServerError");
assert.deepEqual(relayCalls, [true, false]);
assert.equal(
new TextDecoder().decode(stored.get(HOLD_WEBHOOKS_WHILE_OFFLINE_SECRET)),
"false",
);
}),
),
);

it.effect("leaves the relay untouched when the activity setting can't be saved", () =>
withService({ failActivityWrite: true }, ({ preferences, stored, relayCalls }) =>
Effect.gen(function* () {
const error = yield* preferences
.update({ publishAgentActivity: true, holdWebhooksWhileOffline: true })
.pipe(Effect.flip);
assert.equal(error._tag, "EnvironmentHttpInternalServerError");
assert.deepEqual(relayCalls, []);
assert.equal(
new TextDecoder().decode(stored.get(HOLD_WEBHOOKS_WHILE_OFFLINE_SECRET)),
"false",
);
}),
),
);

it.effect("keeps both settings unchanged when the relay refuses the hold change", () =>
withService({ relayFails: true }, ({ preferences, stored }) =>
Effect.gen(function* () {
yield* preferences
.update({ publishAgentActivity: true, holdWebhooksWhileOffline: true })
.pipe(Effect.flip);
assert.equal(new TextDecoder().decode(stored.get(PUBLISH_AGENT_ACTIVITY_SECRET)), "false");
assert.equal(
new TextDecoder().decode(stored.get(HOLD_WEBHOOKS_WHILE_OFFLINE_SECRET)),
"false",
);
}),
),
);

it.effect("changes nothing when the current activity setting can't be read", () =>
withService({ failActivityRead: true, relayFails: true }, ({ preferences, stored, relayCalls }) =>
Effect.gen(function* () {
const error = yield* preferences
.update({ publishAgentActivity: true, holdWebhooksWhileOffline: true })
.pipe(Effect.flip);
assert.equal(error._tag, "EnvironmentHttpInternalServerError");
assert.deepEqual(relayCalls, []);
assert.equal(new TextDecoder().decode(stored.get(PUBLISH_AGENT_ACTIVITY_SECRET)), "false");
}),
),
);

it.effect("two overlapping updates are applied one after the other", () => {
let release!: () => void;
let started!: () => void;
const released = new Promise<void>((resolve) => {
release = resolve;
});
const firstAtRelay = new Promise<void>((resolve) => {
started = resolve;
});
return withService({ holdFirstRelayCall: { started, released } }, ({ preferences, stored }) =>
Effect.gen(function* () {
const first = yield* preferences
.update({ publishAgentActivity: true, holdWebhooksWhileOffline: true })
.pipe(Effect.forkChild);
yield* Effect.promise(() => firstAtRelay);
const second = yield* preferences
.update({ publishAgentActivity: false, holdWebhooksWhileOffline: false })
.pipe(Effect.forkChild);
// Give the second update every chance to run; it must wait for the first.
yield* Effect.yieldNow;
assert.equal(new TextDecoder().decode(stored.get(PUBLISH_AGENT_ACTIVITY_SECRET)), "true");
release();
yield* Fiber.join(first);
yield* Fiber.join(second);
assert.equal(new TextDecoder().decode(stored.get(PUBLISH_AGENT_ACTIVITY_SECRET)), "false");
assert.equal(
new TextDecoder().decode(stored.get(HOLD_WEBHOOKS_WHILE_OFFLINE_SECRET)),
"false",
);
}),
);
});

it.effect("changes nothing when the current hold setting can't be read", () =>
withService({ failHoldRead: true, failHoldWrite: true }, ({ preferences, stored, relayCalls }) =>
Effect.gen(function* () {
stored.set(HOLD_WEBHOOKS_WHILE_OFFLINE_SECRET, encode("true"));
const error = yield* preferences
.update({ publishAgentActivity: true, holdWebhooksWhileOffline: false })
.pipe(Effect.flip);
assert.equal(error._tag, "EnvironmentHttpInternalServerError");
assert.deepEqual(relayCalls, []);
assert.equal(new TextDecoder().decode(stored.get(PUBLISH_AGENT_ACTIVITY_SECRET)), "false");
}),
),
);
130 changes: 130 additions & 0 deletions apps/server/src/cloud/CloudPreferences.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,130 @@
import {
type EnvironmentCloudPreferencesRequest,
EnvironmentHttpBadRequestError,
EnvironmentHttpInternalServerError,
} from "@t3tools/contracts";
import * as Context from "effect/Context";
import * as Effect from "effect/Effect";
import * as Layer from "effect/Layer";
import * as Option from "effect/Option";
import * as Semaphore from "effect/Semaphore";

import * as ServerSecretStore from "../auth/ServerSecretStore.ts";
import * as ServerEnvironment from "../environment/ServerEnvironment.ts";
import * as AgentAwarenessRelay from "../relay/AgentAwarenessRelay.ts";
import { makeRelayEnvironmentClient } from "../relay/relayEnvironmentClient.ts";
import {
HOLD_WEBHOOKS_WHILE_OFFLINE_SECRET,
PUBLISH_AGENT_ACTIVITY_SECRET,
readRelayConnection,
} from "./config.ts";

const encode = (value: boolean) => new TextEncoder().encode(String(value));

/** A failed rollback leaves a setting changed; it is logged, not hidden. */
const rollbackFailed = (cause: unknown) =>
Effect.logWarning("Could not roll back a T3 Connect preference", { cause });

const internalError = (message: string) => (cause: unknown) =>
Effect.logError(message, { cause }).pipe(
Effect.andThen(Effect.fail(new EnvironmentHttpInternalServerError({ message }))),
);

export class CloudPreferences extends Context.Service<
CloudPreferences,
{
/**
* Saves this environment's T3 Connect preferences, all or nothing. The
* activity setting is saved first. Holding webhooks while offline is decided
* by the relay, so the relay is told before the local copy is saved. If
* either step fails, the activity setting is put back, and so is the relay
* when only the local save failed.
*/
readonly update: (
input: EnvironmentCloudPreferencesRequest,
) => Effect.Effect<void, EnvironmentHttpBadRequestError | EnvironmentHttpInternalServerError>;
}
>()("t3/cloud/CloudPreferences") {}

const make = Effect.gen(function* () {
const secrets = yield* ServerSecretStore.ServerSecretStore;
const environment = yield* ServerEnvironment.ServerEnvironment;
const awarenessRelay = yield* AgentAwarenessRelay.AgentAwarenessRelay;
const withSecrets = Effect.provideService(ServerSecretStore.ServerSecretStore, secrets);

const pushHoldWebhooksWhileOffline = Effect.fn("CloudPreferences.pushHoldWebhooksWhileOffline")(
function* (holdWebhooksWhileOffline: boolean) {
const connection = yield* readRelayConnection.pipe(withSecrets);
if (connection === null) {
return yield* new EnvironmentHttpBadRequestError({
message: "Link this environment to T3 Connect first.",
});
}
const environmentId = yield* environment.getEnvironmentId;
const client = yield* makeRelayEnvironmentClient(connection);
yield* client.server
.updateLinkPreferences({
params: { environmentId },
payload: { holdWebhooksWhileOffline },
})
.pipe(
Effect.timeout("10 seconds"),
Effect.catch(internalError("Could not update T3 Connect webhook settings.")),
);
},
);

const save = (name: string, value: boolean) =>
secrets
.set(name, encode(value))
.pipe(Effect.catch(internalError("Could not persist environment cloud preferences.")));

// One update at a time, so two requests can't each leave one setting behind.
const updateLock = yield* Semaphore.make(1);

const update: CloudPreferences["Service"]["update"] = Effect.fn("CloudPreferences.update")(
function* (input) {
Comment thread
macroscopeapp[bot] marked this conversation as resolved.
// All or nothing: the activity setting is saved first, before the relay
// is told anything, and put back if the hold change then fails. Both
// current values are read up front; a failed read stops here, before
// anything changes, because a guessed value would be the rollback target.
const readCurrent = (name: string) =>
secrets
.get(name)
.pipe(Effect.catch(internalError("Could not read environment cloud preferences.")));
const previousActivity = yield* readCurrent(PUBLISH_AGENT_ACTIVITY_SECRET);
const previousHold = yield* readCurrent(HOLD_WEBHOOKS_WHILE_OFFLINE_SECRET);
yield* save(PUBLISH_AGENT_ACTIVITY_SECRET, input.publishAgentActivity);
if (input.holdWebhooksWhileOffline !== undefined) {
const next = input.holdWebhooksWhileOffline;
const previous = Option.match(previousHold, {
onNone: () => false,
onSome: (bytes) => new TextDecoder().decode(bytes) === "true",
});
yield* pushHoldWebhooksWhileOffline(next).pipe(
Effect.andThen(
save(HOLD_WEBHOOKS_WHILE_OFFLINE_SECRET, next).pipe(
Effect.tapError(() =>
previous === next
? Effect.void
: pushHoldWebhooksWhileOffline(previous).pipe(Effect.catch(rollbackFailed)),
),
),
),
Effect.tapError(() =>
Option.match(previousActivity, {
onNone: () => secrets.remove(PUBLISH_AGENT_ACTIVITY_SECRET),
onSome: (bytes) => secrets.set(PUBLISH_AGENT_ACTIVITY_SECRET, bytes),
}).pipe(Effect.catch(rollbackFailed)),
),
);
}
yield* awarenessRelay.requestCatchUp();
Comment thread
macroscopeapp[bot] marked this conversation as resolved.
},
updateLock.withPermits(1),
);

return CloudPreferences.of({ update });
});

export const layer = Layer.effect(CloudPreferences, make);
Loading
Loading