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
Original file line number Diff line number Diff line change
Expand Up @@ -420,6 +420,7 @@ export const makeOrchestrationIntegrationHarness = (
Layer.provideMerge(
Layer.succeed(AgentAwarenessRelay.AgentAwarenessRelay, {
publishThread: () => Effect.void,
requestCatchUp: () => Effect.void,
start: () => Effect.void,
}),
),
Expand Down
9 changes: 9 additions & 0 deletions apps/server/src/cloud/http.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,7 @@ import {
type ServiceUpdateRecord,
} from "./serviceProtocol.ts";
import * as ServerEnvironment from "../environment/ServerEnvironment.ts";
import * as AgentAwarenessRelay from "../relay/AgentAwarenessRelay.ts";
import { CLOUD_CLI_DESIRED_LINK_SECRET } from "./CliState.ts";
import * as CliTokenManager from "./CliTokenManager.ts";
import {
Expand Down Expand Up @@ -81,6 +82,12 @@ const storeFailure = (tag: "AlreadyExists" | "PermissionDenied") =>
});

const unusedSecretStoreOperation = () => Effect.die("unused secret-store operation");
// Linking wakes the awareness relay; these tests do not run it.
const idleAwarenessRelay = AgentAwarenessRelay.AgentAwarenessRelay.of({
publishThread: () => Effect.void,
requestCatchUp: () => Effect.void,
start: () => Effect.void,
});
const decodeManagedTunnelRecoveryRegistration = Schema.decodeUnknownEffect(
Schema.fromJsonString(RelayManagedEndpointRecoveryRegistrationRequest),
);
Expand Down Expand Up @@ -262,6 +269,7 @@ describe("reconcileDesiredCloudLink", () => {
HttpClient.HttpClient,
HttpClient.make(() => unusedSecretStoreOperation()),
),
Effect.provideService(AgentAwarenessRelay.AgentAwarenessRelay, idleAwarenessRelay),
Effect.provide(NodeServices.layer),
),
);
Expand Down Expand Up @@ -375,6 +383,7 @@ describe("releaseManagedTunnelOnShutdown", () => {
<A, E, R>(effect: Effect.Effect<A, E, R>) =>
effect.pipe(
Effect.provideService(ServerSecretStore.ServerSecretStore, harness.store),
Effect.provideService(AgentAwarenessRelay.AgentAwarenessRelay, idleAwarenessRelay),
Effect.provideService(
ServerEnvironment.ServerEnvironment,
ServerEnvironment.ServerEnvironment.of({
Expand Down
5 changes: 5 additions & 0 deletions apps/server/src/cloud/http.ts
Original file line number Diff line number Diff line change
Expand Up @@ -69,6 +69,7 @@ import * as ServerSecretStore from "../auth/ServerSecretStore.ts";
import { requireEnvironmentScope } from "../auth/http.ts";
import * as ServerConfig from "../config.ts";
import * as ServerEnvironment from "../environment/ServerEnvironment.ts";
import * as AgentAwarenessRelay from "../relay/AgentAwarenessRelay.ts";
import * as ManagedEndpointRuntime from "./ManagedEndpointRuntime.ts";
import {
SERVICE_STATE_FILE,
Expand Down Expand Up @@ -411,6 +412,7 @@ interface CloudHttpDependencies {
readonly environmentAuth: EnvironmentAuth.EnvironmentAuth["Service"];
readonly cliTokenManager: CliTokenManager.CloudCliTokenManager["Service"];
readonly httpClient: HttpClient.HttpClient;
readonly awarenessRelay: AgentAwarenessRelay.AgentAwarenessRelay["Service"];
}

const cloudHttpDependencies = Effect.gen(function* () {
Expand All @@ -421,6 +423,7 @@ const cloudHttpDependencies = Effect.gen(function* () {
environmentAuth: yield* EnvironmentAuth.EnvironmentAuth,
cliTokenManager: yield* CliTokenManager.CloudCliTokenManager,
httpClient: yield* HttpClient.HttpClient,
awarenessRelay: yield* AgentAwarenessRelay.AgentAwarenessRelay,
} satisfies CloudHttpDependencies;
});

Expand Down Expand Up @@ -668,6 +671,7 @@ const applyCloudRelayConfig = Effect.fn("environment.cloud.applyRelayConfig")(fu
CLOUD_MINT_PUBLIC_KEY,
stringToBytes(payload.cloudMintPublicKey),
);
yield* dependencies.awarenessRelay.requestCatchUp();
if (payload.endpointRuntime) {
const endpointRuntimeJson = yield* encodeEndpointRuntimeConfigJson(payload.endpointRuntime);
yield* dependencies.secrets.set(
Expand Down Expand Up @@ -1349,6 +1353,7 @@ const cloudPreferencesHandler = Effect.fn("environment.cloud.preferences")(
PUBLISH_AGENT_ACTIVITY_SECRET,
stringToBytes(String(payload.publishAgentActivity)),
);
yield* dependencies.awarenessRelay.requestCatchUp();
Comment thread
macroscopeapp[bot] marked this conversation as resolved.
return yield* readCloudLinkState(dependencies);
},
Effect.catchIf(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -108,6 +108,7 @@ describe("OrchestrationReactor", () => {
Layer.provideMerge(
Layer.succeed(AgentAwarenessRelay.AgentAwarenessRelay, {
publishThread: () => Effect.void,
requestCatchUp: () => Effect.void,
start: () => {
started.push("agent-awareness-relay");
return Effect.void;
Expand Down
18 changes: 11 additions & 7 deletions apps/server/src/persistence/ProviderSessionRuntime.ts
Original file line number Diff line number Diff line change
Expand Up @@ -102,11 +102,14 @@ export class ProviderSessionRuntimeRepository extends Context.Service<
>;

/**
* List all provider runtime rows.
* List provider runtime rows.
*
* Returned in ascending last-seen order.
* Returned in ascending last-seen order. `excludeStopped` filters stopped
* rows in SQL. Long-lived installs keep thousands for their resume cursors.
*/
readonly list: () => Effect.Effect<
readonly list: (options?: {
readonly excludeStopped?: boolean;
}) => Effect.Effect<
ReadonlyArray<ProviderSessionRuntime>,
ProviderSessionRuntimeRepositoryError
>;
Expand Down Expand Up @@ -336,9 +339,9 @@ export const make = Effect.gen(function* () {
});

const listRuntimeRows = SqlSchema.findAll({
Request: Schema.Void,
Request: Schema.Struct({ excludeStopped: Schema.Boolean }),
Result: ProviderSessionRuntimeRawDbRowSchema,
execute: () =>
execute: ({ excludeStopped }) =>
sql`
SELECT
thread_id AS "threadId",
Expand All @@ -351,6 +354,7 @@ export const make = Effect.gen(function* () {
resume_cursor_json AS "resumeCursor",
runtime_payload_json AS "runtimePayload"
FROM provider_session_runtime
${excludeStopped ? sql`WHERE status != 'stopped'` : sql``}
ORDER BY last_seen_at ASC, thread_id ASC
`,
});
Expand Down Expand Up @@ -414,8 +418,8 @@ export const make = Effect.gen(function* () {
),
);

const list: ProviderSessionRuntimeRepository["Service"]["list"] = () =>
listRuntimeRows(undefined).pipe(
const list: ProviderSessionRuntimeRepository["Service"]["list"] = (options) =>
listRuntimeRows({ excludeStopped: options?.excludeStopped === true }).pipe(
Effect.mapError(
toPersistenceSqlOrDecodeError(
"ProviderSessionRuntimeRepository.list:query",
Expand Down
35 changes: 35 additions & 0 deletions apps/server/src/provider/Layers/ProviderSessionDirectory.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -372,6 +372,41 @@ it.layer(makeDirectoryLayer(SqlitePersistenceMemory))("ProviderSessionDirectoryL
}),
);

it.effect("lists only bindings that are not stopped when asked", () =>
Effect.gen(function* () {
const directory = yield* ProviderSessionDirectory;
const runtimeRepository = yield* ProviderSessionRuntime.ProviderSessionRuntimeRepository;
const statuses = ["running", "starting", "error", "stopped"] as const;
const threadIds = new Set<string>();

for (const status of statuses) {
const threadId = ThreadId.make(`thread-exclude-stopped-${status}`);
threadIds.add(threadId);
yield* runtimeRepository.upsert({
threadId,
providerName: "codex",
providerInstanceId: ProviderInstanceId.make("codex"),
adapterKey: "codex",
runtimeMode: "full-access",
status,
lastSeenAt: "2026-04-14T12:00:00.000Z",
resumeCursor: null,
runtimePayload: null,
});
}

const liveStatuses = (yield* directory.listBindings({ excludeStopped: true }))
.filter((binding) => threadIds.has(binding.threadId))
.map((binding) => binding.status);
const allStatuses = (yield* directory.listBindings())
.filter((binding) => threadIds.has(binding.threadId))
.map((binding) => binding.status);

assert.deepEqual(liveStatuses.toSorted(), ["error", "running", "starting"]);
assert.deepEqual(allStatuses.toSorted(), ["error", "running", "starting", "stopped"]);
}),
);

it.effect(
"resets adapterKey to the new provider when provider changes without an explicit adapter key",
() =>
Expand Down
4 changes: 2 additions & 2 deletions apps/server/src/provider/Layers/ProviderSessionDirectory.ts
Original file line number Diff line number Diff line change
Expand Up @@ -182,8 +182,8 @@ const makeProviderSessionDirectory = Effect.gen(function* () {
Effect.map((rows) => rows.map((row) => row.threadId)),
);

const listBindings: ProviderSessionDirectoryShape["listBindings"] = () =>
repository.list().pipe(
const listBindings: ProviderSessionDirectoryShape["listBindings"] = (options) =>
repository.list(options).pipe(
Effect.mapError(toPersistenceError("ProviderSessionDirectory.listBindings:list")),
Effect.flatMap((rows) =>
Effect.forEach(
Expand Down
10 changes: 4 additions & 6 deletions apps/server/src/provider/Layers/ProviderSessionReaper.ts
Original file line number Diff line number Diff line change
Expand Up @@ -35,15 +35,13 @@ const makeProviderSessionReaper = (options?: ProviderSessionReaperLiveOptions) =
const sweepIntervalMs = Math.max(1, options?.sweepIntervalMs ?? DEFAULT_SWEEP_INTERVAL_MS);

const sweep = Effect.gen(function* () {
const bindings = yield* directory.listBindings();
// Stopped rows stay for their resume cursors and far outnumber live
// ones, so the query skips them.
const bindings = yield* directory.listBindings({ excludeStopped: true });
const now = yield* Clock.currentTimeMillis;
let reapedCount = 0;

for (const binding of bindings) {
if (binding.status === "stopped") {
continue;
}

const lastSeenMs = Date.parse(binding.lastSeenAt);
if (Number.isNaN(lastSeenMs)) {
yield* Effect.logWarning("provider.session.reaper.invalid-last-seen", {
Expand Down Expand Up @@ -122,7 +120,7 @@ const makeProviderSessionReaper = (options?: ProviderSessionReaperLiveOptions) =
if (reapedCount > 0) {
yield* Effect.logInfo("provider.session.reaper.sweep-complete", {
reapedCount,
totalBindings: bindings.length,
liveBindings: bindings.length,
});
}
});
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -70,7 +70,10 @@ export interface ProviderSessionDirectoryShape {
ProviderSessionDirectoryPersistenceError
>;

readonly listBindings: () => Effect.Effect<
/** `excludeStopped` skips stopped rows in the query, not after decoding. */
readonly listBindings: (options?: {
readonly excludeStopped?: boolean;
}) => Effect.Effect<
ReadonlyArray<ProviderRuntimeBindingWithMetadata>,
ProviderSessionDirectoryPersistenceError
>;
Expand Down
112 changes: 112 additions & 0 deletions apps/server/src/relay/AgentAwarenessRelay.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,7 @@ import * as Option from "effect/Option";
import * as Queue from "effect/Queue";
import * as Stream from "effect/Stream";
import * as Tracer from "effect/Tracer";
import * as TestClock from "effect/testing/TestClock";

import * as ServerSecretStore from "../auth/ServerSecretStore.ts";
import * as ServerEnvironment from "../environment/ServerEnvironment.ts";
Expand Down Expand Up @@ -975,3 +976,114 @@ describe.sequential("signRelayAgentActivityPublishProof", () => {
),
);
});

describe.sequential("startup catch-up", () => {
// An unlinked relay with publishing off. `link` writes the link secrets and
// `enablePublishing` the opt-in. Counts link checks (relay URL reads) and
// catch-up publishes (shell snapshot reads).
function makeUnlinkedRelay() {
const secrets = makeMemorySecretStore();
const counts = { linkChecks: 0, catchUpPublishes: 0 };
const countingStore = {
...secrets.store,
get: (name: string) =>
Effect.suspend(() => {
if (name === RELAY_URL_SECRET) counts.linkChecks += 1;
return secrets.store.get(name);
}),
} satisfies ServerSecretStore.ServerSecretStore["Service"];

const layer = AgentAwarenessRelay.layer.pipe(
Layer.provide(
Layer.mergeAll(
Layer.succeed(ServerSecretStore.ServerSecretStore, countingStore),
Layer.succeed(ServerEnvironment.ServerEnvironment, {
getEnvironmentId: Effect.succeed("env-1" as EnvironmentId),
getDescriptor: Effect.die("unused descriptor"),
}),
Layer.succeed(OrchestrationEngineService, {
streamDomainEvents: Stream.never,
} as unknown as OrchestrationEngineShape),
Layer.succeed(ProjectionSnapshotQuery, {
getShellSnapshot: () =>
Effect.sync(() => {
counts.catchUpPublishes += 1;
return {
snapshotSequence: 1,
projects: [],
threads: [],
updatedAt: "2026-05-25T00:00:00.000Z",
} satisfies OrchestrationShellSnapshot;
}),
} as unknown as ProjectionSnapshotQueryShape),
),
),
Layer.provideMerge(NodeServices.layer),
);
const link = Effect.all(
[
secrets.setString(RELAY_URL_SECRET, "https://relay.example.test"),
secrets.setString(RELAY_ENVIRONMENT_CREDENTIAL_SECRET, "relay-credential"),
],
{ discard: true },
);
const enablePublishing = secrets.setString(PUBLISH_AGENT_ACTIVITY_SECRET, "true");
return { counts, layer, link, enablePublishing };
}

it.effect("checks an unlinked environment once a minute and still catches up once linked", () => {
const { counts, layer, link, enablePublishing } = makeUnlinkedRelay();
return Effect.gen(function* () {
const relay = yield* AgentAwarenessRelay.AgentAwarenessRelay;
yield* enablePublishing;
yield* relay.start();

// Get past the backoff ramp, then count checks in a steady window.
yield* TestClock.adjust("10 minutes");
const checksBeforeWindow = counts.linkChecks;
yield* TestClock.adjust("10 minutes");
expect(counts.linkChecks - checksBeforeWindow).toBe(10);
expect(counts.catchUpPublishes).toBe(0);

yield* link;
yield* TestClock.adjust("1 minute");
expect(counts.catchUpPublishes).toBe(1);
}).pipe(Effect.provide(layer), Effect.scoped);
});

it.effect("publishes at once when this process links while the check is backed off", () => {
const { counts, layer, link, enablePublishing } = makeUnlinkedRelay();
return Effect.gen(function* () {
const relay = yield* AgentAwarenessRelay.AgentAwarenessRelay;
yield* enablePublishing;
yield* relay.start();

// Backed off to 60 s: the next check is still seconds away.
yield* TestClock.adjust("10 minutes");
yield* link;
yield* TestClock.adjust("1 second");
expect(counts.catchUpPublishes).toBe(0);

yield* relay.requestCatchUp();
yield* TestClock.adjust("1 second");
expect(counts.catchUpPublishes).toBe(1);
}).pipe(Effect.provide(layer), Effect.scoped);
});

it.effect("catches up within 5 s when another process enables publishing on a link", () => {
const { counts, layer, link, enablePublishing } = makeUnlinkedRelay();
return Effect.gen(function* () {
const relay = yield* AgentAwarenessRelay.AgentAwarenessRelay;
yield* link;
yield* relay.start();

yield* TestClock.adjust("10 minutes");
expect(counts.catchUpPublishes).toBe(0);

// `t3 connect publish` writes the opt-in without waking this process.
yield* enablePublishing;
yield* TestClock.adjust("5 seconds");
expect(counts.catchUpPublishes).toBe(1);
}).pipe(Effect.provide(layer), Effect.scoped);
});
});
Loading
Loading