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
3 changes: 2 additions & 1 deletion apps/server/src/server.ts
Original file line number Diff line number Diff line change
Expand Up @@ -49,7 +49,7 @@ import { ProviderRuntimeIngestionLive } from "./orchestration/Layers/ProviderRun
import { ProviderCommandReactorLive } from "./orchestration/Layers/ProviderCommandReactor.ts";
import { CheckpointReactorLive } from "./orchestration/Layers/CheckpointReactor.ts";
import { ThreadDeletionReactorLive } from "./orchestration/Layers/ThreadDeletionReactor.ts";
import { T3xLayerLive } from "./t3x/index.ts"; // t3x fork seam — all fork-local features fan in here
import { T3xLayerLive, T3xRoutesLive } from "./t3x/index.ts"; // t3x fork seam — all fork-local features fan in here
import * as AgentAwarenessRelay from "./relay/AgentAwarenessRelay.ts";
import { hasCloudPublicConfig } from "./cloud/publicConfig.ts";
import { ProviderRegistryLive } from "./provider/Layers/ProviderRegistry.ts";
Expand Down Expand Up @@ -361,6 +361,7 @@ export const makeRoutesLayer = Layer.mergeAll(
Layer.provide(environmentAuthenticatedAuthLayer),
),
otlpTracesProxyRouteLayer,
T3xRoutesLive, // t3x fork seam — all fork-local HTTP routes fan in here
assetRouteLayer,
staticAndDevRouteLayer,
websocketRpcRouteLayer,
Expand Down
47 changes: 47 additions & 0 deletions apps/server/src/t3x/autoResume/Reactor.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -255,4 +255,51 @@ describe("AutoResumeReactor (integration)", () => {
}).pipe(Effect.provide(AutoResumeReactorLive.pipe(Layer.provideMerge(deps))));
}).pipe(Effect.scoped, Effect.provide(Layer.mergeAll(NodeServices.layer, TestClock.layer()))),
);

it.effect("does NOT schedule for a thread whose auto-resume is switched off", () =>
Effect.gen(function* () {
const { dispatched, deps, store } = yield* harness(readModel({}), [rejectedEvent(100)]);
yield* store.setEnabled("thread-1", false);

yield* Effect.gen(function* () {
yield* settle; // detection runs against the pre-loaded rejection

assert.strictEqual(
(yield* store.listPending).length,
0,
"a disabled thread must never schedule a resume",
);
// Disabling is a deliberate user action, so it must not post timeline noise either.
assert.notInclude(types(yield* Ref.get(dispatched)), "thread.activity.append");

yield* advancePastResume;
assert.notInclude(types(yield* Ref.get(dispatched)), "thread.turn.start");
}).pipe(Effect.provide(AutoResumeReactorLive.pipe(Layer.provideMerge(deps))));
}).pipe(Effect.scoped, Effect.provide(Layer.mergeAll(NodeServices.layer, TestClock.layer()))),
);

it.effect("cancels an already-scheduled resume when the thread is switched off mid-wait", () =>
Effect.gen(function* () {
const { dispatched, deps, store } = yield* harness(readModel({}), [rejectedEvent(100)]);

yield* Effect.gen(function* () {
yield* settle;
assert.strictEqual((yield* store.listPending).length, 1, "precondition: it scheduled");

// The switch is flipped off *after* scheduling but *before* the window reopens.
// fireOne must re-read the record rather than trust the scheduling-time value.
yield* store.setEnabled("thread-1", false);

yield* advancePastResume;

assert.notInclude(types(yield* Ref.get(dispatched)), "thread.turn.start");
assert.strictEqual((yield* store.listPending).length, 0, "pending must be cleared");
assert.strictEqual(
yield* store.countFiredSince("thread-1", 0),
0,
"a cancellation must not burn one of the 24h attempts",
);
}).pipe(Effect.provide(AutoResumeReactorLive.pipe(Layer.provideMerge(deps))));
}).pipe(Effect.scoped, Effect.provide(Layer.mergeAll(NodeServices.layer, TestClock.layer()))),
);
});
18 changes: 18 additions & 0 deletions apps/server/src/t3x/autoResume/Reactor.ts
Original file line number Diff line number Diff line change
Expand Up @@ -133,6 +133,9 @@ const makeSupervisor = Effect.gen(function* () {

const threadId = event.threadId;
const record = yield* store.getThread(threadId);
// Per-thread switch. A thread the user turned off never schedules, and posts no
// timeline note: they disabled it deliberately, so an activity row would be noise.
if (!record.enabled) return;
const nowMs = yield* Clock.currentTimeMillis;
const firedRecently = yield* store.countFiredSince(threadId, nowMs - BACKOFF_LOOKBACK_MS);
const firedInCapWindow = yield* store.countFiredSince(threadId, nowMs - CAP_WINDOW_MS);
Expand Down Expand Up @@ -183,6 +186,21 @@ const makeSupervisor = Effect.gen(function* () {
return;
}

// The user can switch a thread off *after* its resume was scheduled, so re-read the
// record here rather than trusting the value seen at scheduling time, and treat a
// disabled thread as a cancellation (drop the pending resume, say so once).
const record = yield* store.getThread(pending.threadId);
if (!record.enabled) {
yield* store.clearPending(pending.threadId);
yield* appendActivity(
thread.id,
"info",
"t3x.auto-resume.cancelled",
"Auto-resume cancelled: turned off for this thread.",
);
return;
}

const reason = cancelReason(thread, pending.baseline);
if (reason) {
yield* store.clearPending(pending.threadId);
Expand Down
216 changes: 216 additions & 0 deletions apps/server/src/t3x/autoResume/http.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,216 @@
// @effect-diagnostics nodeBuiltinImport:off
/**
* Route-level tests for `/api/t3x/auto-resume`.
*
* Serves ONLY the fork's route layer (not the whole `makeRoutesLayer`) over a real HTTP
* server on an ephemeral port, with `EnvironmentAuth` mocked.
*
* Scope note: because auth is mocked, these tests do NOT prove that the route's
* `authenticateWithOperateScope` faithfully mirrors upstream's private
* `authenticateRawRouteWithScope` — that remains a logic mirror tracked in SEAMS.md.
* What they DO prove is that a rejected credential and a missing scope each render as a
* 401/403 through the route's `catchTags` rather than escaping as a 500 or an unhandled
* defect, and that the GET/POST contract round-trips through a real store.
*/
import * as NodeHttpServer from "@effect/platform-node/NodeHttpServer";
import * as NodeServices from "@effect/platform-node/NodeServices";
import { assert, describe, it } from "@effect/vitest";
import { AuthOrchestrationOperateScope } from "@t3tools/contracts";
import * as Effect from "effect/Effect";
import * as FileSystem from "effect/FileSystem";
import * as Layer from "effect/Layer";
import {
FetchHttpClient,
HttpClient,
HttpClientRequest,
HttpRouter,
HttpServer,
} from "effect/unstable/http";
import * as NodeHttp from "node:http";
import * as NodePath from "node:path";

import * as EnvironmentAuth from "../../auth/EnvironmentAuth.ts";
import { autoResumeRouteLayer } from "./http.ts";
import { AutoResumeStore, makeAutoResumeStore } from "./state.ts";

const PATH = "/api/t3x/auto-resume";

const authedSession = (scopes: ReadonlyArray<string>) =>
({
sessionId: "test-session",
subject: "test",
method: "bearer-access-token",
scopes,
}) as unknown as EnvironmentAuth.AuthenticatedSession;

/** Mock auth that succeeds with the given scopes. */
const authOk = (scopes: ReadonlyArray<string> = [AuthOrchestrationOperateScope]) =>
Layer.mock(EnvironmentAuth.EnvironmentAuth)({
authenticateHttpRequest: () => Effect.succeed(authedSession(scopes)),
});

/**
* Mock auth that rejects the credential.
*
* Must be a member of the `ServerAuthCredentialError` union
* (`ServerAuthMissingCredentialError | ServerAuthInvalidCredentialError`) — that union is
* what `isServerAuthCredentialError` guards on. An error outside it (and outside
* `ServerAuthInternalError`) matches neither `catchIf` branch and surfaces as a 500. That
* is equally true of upstream's OTLP route, which uses the identical two branches, so it
* is a shared property rather than a fork divergence.
*/
const authRejects = Layer.mock(EnvironmentAuth.EnvironmentAuth)({
authenticateHttpRequest: () =>
Effect.fail(new EnvironmentAuth.ServerAuthInvalidCredentialError({})),
});

const baseUrl = (pathname: string) =>
Effect.gen(function* () {
const server = yield* HttpServer.HttpServer;
const address = server.address as HttpServer.TcpAddress;
return `http://127.0.0.1:${address.port}${pathname}`;
});

const getJson = (pathname: string) =>
Effect.gen(function* () {
const client = yield* HttpClient.HttpClient;
return yield* client.get(yield* baseUrl(pathname));
});

const postJson = (pathname: string, body: unknown) =>
Effect.gen(function* () {
const client = yield* HttpClient.HttpClient;
const request = HttpClientRequest.bodyJsonUnsafe(
HttpClientRequest.post(yield* baseUrl(pathname)),
body,
);
return yield* client.execute(request);
});

/** Serves the route over an ephemeral port with the given auth layer and a real store. */
const withRoute = <A, E>(
authLayer: Layer.Layer<EnvironmentAuth.EnvironmentAuth>,
body: Effect.Effect<A, E, HttpServer.HttpServer | HttpClient.HttpClient>,
): Promise<A> =>
Effect.gen(function* () {
const fs = yield* FileSystem.FileSystem;
const root = yield* fs.makeTempDirectoryScoped({ prefix: "t3x-http-" });
const storeLayer = Layer.effect(
AutoResumeStore,
makeAutoResumeStore(NodePath.join(root, "state.json")),
);

const served = HttpRouter.serve(autoResumeRouteLayer, {
disableListenLog: true,
disableLogger: true,
}).pipe(
Layer.provide(storeLayer),
Layer.provide(authLayer),
// provideMerge, not provide: the test body needs HttpServer in context to read the
// ephemeral port off `server.address`.
Layer.provideMerge(NodeHttpServer.layer(() => NodeHttp.createServer(), { port: 0 })),
);

return yield* body.pipe(Effect.provide(served), Effect.provide(FetchHttpClient.layer));
}).pipe(Effect.scoped, Effect.provide(NodeServices.layer), Effect.runPromise);

describe("/api/t3x/auto-resume", () => {
it("rejects a bad credential with 401/403 instead of leaking a 500", () =>
withRoute(
authRejects,
Effect.gen(function* () {
const res = yield* getJson(`${PATH}?threadId=thread-a`);
assert.isTrue(
res.status === 401 || res.status === 403,
`expected 401/403 for a rejected credential, got ${res.status}`,
);
}),
));

it("rejects a session lacking the operate scope", () =>
withRoute(
authOk([]),
Effect.gen(function* () {
const res = yield* getJson(`${PATH}?threadId=thread-a`);
assert.isTrue(
res.status === 401 || res.status === 403,
`expected 401/403 without the operate scope, got ${res.status}`,
);
}),
));

it("GET defaults a never-seen thread to enabled", () =>
withRoute(
authOk(),
Effect.gen(function* () {
const res = yield* getJson(`${PATH}?threadId=fresh`);
assert.strictEqual(res.status, 200);
const body = (yield* res.json) as Record<string, unknown>;
assert.strictEqual(body.enabled, true);
assert.strictEqual(body.overridePrompt, null);
assert.strictEqual(body.pending, null);
}),
));

it("GET without a threadId is a 400, not a crash", () =>
withRoute(
authOk(),
Effect.gen(function* () {
const res = yield* getJson(PATH);
assert.strictEqual(res.status, 400);
}),
));

it("POST patches one field without clobbering the other", () =>
withRoute(
authOk(),
Effect.gen(function* () {
const off = yield* postJson(PATH, { threadId: "t1", enabled: false });
assert.strictEqual(off.status, 200);
assert.strictEqual(((yield* off.json) as Record<string, unknown>).enabled, false);

// Writing only overridePrompt must not resurrect `enabled`.
const text = yield* postJson(PATH, { threadId: "t1", overridePrompt: "keep going" });
const textBody = (yield* text.json) as Record<string, unknown>;
assert.strictEqual(textBody.overridePrompt, "keep going");
assert.strictEqual(textBody.enabled, false, "a prompt write must not clobber enabled");

// …and writing only `enabled` must not wipe the prompt.
const on = yield* postJson(PATH, { threadId: "t1", enabled: true });
const onBody = (yield* on.json) as Record<string, unknown>;
assert.strictEqual(onBody.enabled, true);
assert.strictEqual(onBody.overridePrompt, "keep going", "toggle must not clobber prompt");
}),
));

// This route is what makes threadId caller-controlled; before it, ids only ever came
// from provider events. A prototype-chain id used to resolve on Object.prototype, which
// 500'd on read and persisted a malformed record that wiped the whole state file on the
// next boot. Pinned here at the HTTP boundary, not just in the store.
it("handles a prototype-chain threadId as an ordinary unknown thread", () =>
withRoute(
authOk(),
Effect.gen(function* () {
for (const hostile of ["constructor", "__proto__", "toString"]) {
const res = yield* getJson(`${PATH}?threadId=${encodeURIComponent(hostile)}`);
assert.strictEqual(res.status, 200, `GET threadId=${hostile} must not 500`);
const body = (yield* res.json) as Record<string, unknown>;
assert.strictEqual(body.enabled, true);
assert.strictEqual(body.pending, null);
}

const wrote = yield* postJson(PATH, { threadId: "__proto__", enabled: false });
assert.strictEqual(wrote.status, 200, "POST threadId=__proto__ must not 500");
assert.strictEqual(((yield* wrote.json) as Record<string, unknown>).enabled, false);
}),
));

it("POST with a malformed body is a 400, not a 500", () =>
withRoute(
authOk(),
Effect.gen(function* () {
const res = yield* postJson(PATH, { nope: true });
assert.strictEqual(res.status, 400);
}),
));
});
Loading