Skip to content
Merged
Show file tree
Hide file tree
Changes from 3 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
1 change: 1 addition & 0 deletions apps/desktop/src/settings/DesktopClientSettings.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,7 @@ const clientSettings: ClientSettings = {
confirmThreadArchive: true,
confirmThreadDelete: false,
confirmThreadUnpin: false,
continueThreadsAfterServerUpdate: true,
dismissedProviderUpdateNotificationKeys: [],
diffIgnoreWhitespace: true,
environmentIdentificationMode: "artwork",
Expand Down
105 changes: 105 additions & 0 deletions apps/server/src/cloud/selfUpdate.test.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
import * as NodeServices from "@effect/platform-node/NodeServices";
import { expect, it } from "@effect/vitest";
import { ServerSelfUpdateError, ThreadId } from "@t3tools/contracts";
import { HostProcessExecutablePath } from "@t3tools/shared/hostProcess";
import * as Deferred from "effect/Deferred";
import * as Effect from "effect/Effect";
Expand Down Expand Up @@ -104,6 +105,110 @@ const makeHarness = Effect.fn("test.make_self_update_harness")(function* (
});

it.layer(NodeServices.layer)("server self update", (it) => {
it.effect("marks running threads at the boot-service handoff", () =>
Effect.gen(function* () {
const events: string[] = [];
const selfUpdate = yield* ServerSelfUpdate.withRunningThreadContinuation({
mode: "web",
selfUpdate: {
update: (_input, reportProgress = () => Effect.void) =>
reportProgress("downloading").pipe(
Effect.andThen(reportProgress("installing")),
Effect.as({
targetVersion: "1.1.0",
method: "boot-service" as const,
updateId: "update-id",
}),
),
commitDesktopUpdate: () => Effect.never,
},
prepare: Effect.sync(() => {
events.push("prepare");
return [ThreadId.make("thread-running")];
}),
clear: () => Effect.sync(() => void events.push("clear")),
});

yield* selfUpdate.update({ targetVersion: "1.1.0", continueRunningThreads: true }, (stage) =>
Effect.sync(() => void events.push(stage)),
);

expect(events).toEqual(["downloading", "prepare", "installing"]);
}),
);

it.effect("marks desktop threads only when the prepared update commits", () =>
Effect.gen(function* () {
const threadId = ThreadId.make("thread-running-desktop");
const events: string[] = [];
const commitError = new ServerSelfUpdateError({ reason: "install failed" });
const selfUpdate = yield* ServerSelfUpdate.withRunningThreadContinuation({
mode: "desktop",
selfUpdate: {
update: (_input, reportProgress = () => Effect.void) =>
reportProgress("installing").pipe(
Effect.as({
targetVersion: "1.2.0",
method: "desktop-app" as const,
desktopUpdateToken: "desktop-token",
}),
),
commitDesktopUpdate: () =>
Effect.sync(() => events.push("commit")).pipe(Effect.andThen(Effect.fail(commitError))),
},
prepare: Effect.sync(() => {
events.push("prepare");
return [threadId];
}),
clear: (threadIds) => Effect.sync(() => void events.push(`clear:${threadIds.join(",")}`)),
});

yield* selfUpdate.update({ targetVersion: "1.2.0", continueRunningThreads: true }, (stage) =>
Effect.sync(() => void events.push(stage)),
);
expect(events).toEqual(["installing"]);
expect(yield* selfUpdate.commitDesktopUpdate("desktop-token").pipe(Effect.flip)).toBe(
commitError,
);
expect(events).toEqual(["installing", "prepare", "commit", `clear:${threadId}`]);
expect(yield* selfUpdate.commitDesktopUpdate("desktop-token").pipe(Effect.flip)).toBe(
commitError,
);
expect(events).toEqual([
"installing",
"prepare",
"commit",
`clear:${threadId}`,
"prepare",
"commit",
`clear:${threadId}`,
]);
}),
);

it.effect("reports a failed continuation-marker cleanup", () =>
Effect.gen(function* () {
const updateError = new ServerSelfUpdateError({ reason: "update failed" });
const clearError = new ServerSelfUpdateError({ reason: "marker cleanup failed" });
const selfUpdate = yield* ServerSelfUpdate.withRunningThreadContinuation({
mode: "web",
selfUpdate: {
update: (_input, reportProgress = () => Effect.void) =>
reportProgress("installing").pipe(Effect.andThen(Effect.fail(updateError))),
commitDesktopUpdate: () => Effect.never,
},
prepare: Effect.succeed([ThreadId.make("thread-cleanup-failure")]),
clear: () => Effect.fail(clearError),
});

expect(
yield* selfUpdate
.update({ targetVersion: "1.1.0", continueRunningThreads: true })
.pipe(Effect.flip),
).toBe(clearError);
}),
);

it.effect("stages and preflights before asking the launcher for an update ID", () =>
Effect.gen(function* () {
const { selfUpdate, order } = yield* makeHarness();
Expand Down
97 changes: 95 additions & 2 deletions apps/server/src/cloud/selfUpdate.ts
Original file line number Diff line number Diff line change
Expand Up @@ -4,15 +4,17 @@ import {
type ServerSelfUpdateInput,
type ServerSelfUpdateProgressStage,
type ServerSelfUpdateResult,
type ThreadId,
} from "@t3tools/contracts";
import { HostProcessExecutablePath } from "@t3tools/shared/hostProcess";
import * as Context from "effect/Context";
import * as Duration from "effect/Duration";
import * as Effect from "effect/Effect";
import * as HashSet from "effect/HashSet";
import * as Ref from "effect/Ref";
import * as FileSystem from "effect/FileSystem";
import * as Layer from "effect/Layer";
import * as Path from "effect/Path";
import * as Ref from "effect/Ref";

import * as ServerConfig from "../config.ts";
import * as DesktopAppUpdate from "../desktopUpdate/DesktopAppUpdate.ts";
Expand Down Expand Up @@ -41,14 +43,105 @@ export class ServerSelfUpdate extends Context.Service<
{
readonly update: (
input: ServerSelfUpdateInput,
reportProgress?: (stage: ServerSelfUpdateProgressStage) => Effect.Effect<void>,
reportProgress?: (
stage: ServerSelfUpdateProgressStage,
) => Effect.Effect<void, ServerSelfUpdateError>,
) => Effect.Effect<ServerSelfUpdateResult, ServerSelfUpdateError>;
readonly commitDesktopUpdate: (
requestId: string,
) => Effect.Effect<never, ServerSelfUpdateError>;
}
>()("t3/cloud/selfUpdate/ServerSelfUpdate") {}

export const withRunningThreadContinuation = Effect.fn(
"cloud.server_self_update.withRunningThreadContinuation",
)(function* (input: {
readonly mode: ServerConfig.RuntimeMode;
readonly selfUpdate: ServerSelfUpdate["Service"];
readonly prepare: Effect.Effect<ReadonlyArray<ThreadId>, ServerSelfUpdateError>;
readonly clear: (
threadIds: ReadonlyArray<ThreadId>,
) => Effect.Effect<void, ServerSelfUpdateError>;
}) {
const desktopContinuationTokens = yield* Ref.make(HashSet.empty<string>());
const clearOnError = <A>(
effect: Effect.Effect<A, ServerSelfUpdateError>,
threadIds: () => ReadonlyArray<ThreadId>,
): Effect.Effect<A, ServerSelfUpdateError> =>
effect.pipe(
Effect.catchCause((cause) =>
Comment thread
macroscopeapp[bot] marked this conversation as resolved.
input.clear(threadIds()).pipe(Effect.andThen(Effect.failCause(cause))),
),
);
Comment thread
cursor[bot] marked this conversation as resolved.

const update: ServerSelfUpdate["Service"]["update"] = (
request,
reportProgress = () => Effect.void,
) => {
let prepared = false;
let continuationThreadIds: ReadonlyArray<ThreadId> = [];
return clearOnError(
input.selfUpdate
.update(request, (stage) =>
(request.continueRunningThreads === true &&
input.mode !== "desktop" &&
stage === "installing" &&
!prepared
? input.prepare.pipe(
Effect.tap((threadIds) =>
Effect.sync(() => {
prepared = true;
continuationThreadIds = threadIds;
}),
),
Effect.asVoid,
)
: Effect.void
).pipe(Effect.andThen(reportProgress(stage))),
)
.pipe(
Effect.tap((result) => {
if (
result.method === "desktop-app" &&
result.desktopUpdateToken !== undefined &&
request.continueRunningThreads === true
) {
return Ref.update(desktopContinuationTokens, HashSet.add(result.desktopUpdateToken));
}
return Effect.void;
}),
),
() => continuationThreadIds,
);
};

return ServerSelfUpdate.of({
update,
commitDesktopUpdate: (requestId) =>
Effect.gen(function* () {
const shouldContinue = yield* Ref.modify(desktopContinuationTokens, (tokens) => [
HashSet.has(tokens, requestId),
HashSet.remove(tokens, requestId),
]);
let continuationThreadIds: ReadonlyArray<ThreadId> = [];
return yield* clearOnError(
Effect.gen(function* () {
continuationThreadIds = shouldContinue ? yield* input.prepare : [];
return yield* input.selfUpdate.commitDesktopUpdate(requestId);
}),
() => continuationThreadIds,
).pipe(
Effect.catchCause((cause) =>
(shouldContinue
? Ref.update(desktopContinuationTokens, HashSet.add(requestId))
: Effect.void
).pipe(Effect.andThen(Effect.failCause(cause))),
),
);
}),
});
});

export const make = Effect.fn("cloud.server_self_update.make")(function* () {
const serverConfig = yield* ServerConfig.ServerConfig;
const desktopAppUpdate = yield* DesktopAppUpdate.DesktopAppUpdate;
Expand Down
16 changes: 12 additions & 4 deletions apps/server/src/desktopUpdate/DesktopAppUpdate.ts
Original file line number Diff line number Diff line change
Expand Up @@ -49,7 +49,9 @@ export class DesktopAppUpdate extends Context.Service<
/** Checks and downloads through the desktop app, then returns a token
while this server is still connected. `commit` starts installation. */
readonly run: (
reportProgress: (stage: ServerSelfUpdateProgressStage) => Effect.Effect<void>,
reportProgress: (
stage: ServerSelfUpdateProgressStage,
) => Effect.Effect<void, ServerSelfUpdateError>,
) => Effect.Effect<ServerSelfUpdateResult, ServerSelfUpdateError>;
/** Starts the prepared install. Success stops this server, so this effect
returns only when installation fails or times out. */
Expand All @@ -72,11 +74,15 @@ export const make = Effect.fn("desktopUpdate.desktopAppUpdate.make")(function* (
const consumeReports = (
requestId: string,
changes: Stream.Stream<DesktopUpdateStatusReport>,
reportProgress: (stage: ServerSelfUpdateProgressStage) => Effect.Effect<void>,
reportProgress: (
stage: ServerSelfUpdateProgressStage,
) => Effect.Effect<void, ServerSelfUpdateError>,
) =>
Effect.gen(function* () {
const lastStage = yield* Ref.make<ServerSelfUpdateProgressStage | null>(null);
const emitStage = (stage: ServerSelfUpdateProgressStage | null): Effect.Effect<void> =>
const emitStage = (
stage: ServerSelfUpdateProgressStage | null,
): Effect.Effect<void, ServerSelfUpdateError> =>
stage === null
? Effect.void
: Ref.get(lastStage).pipe(
Expand All @@ -90,7 +96,9 @@ export const make = Effect.fn("desktopUpdate.desktopAppUpdate.make")(function* (
const terminal = yield* changes.pipe(
Stream.filter((report) => report.requestId === requestId),
Stream.mapEffect(
(report): Effect.Effect<Option.Option<DesktopUpdateStatusReport>> =>
(
report,
): Effect.Effect<Option.Option<DesktopUpdateStatusReport>, ServerSelfUpdateError> =>
report.outcome === undefined
? emitStage(desktopUpdateProgressStage(report.state)).pipe(
Effect.as(Option.none<DesktopUpdateStatusReport>()),
Expand Down
2 changes: 2 additions & 0 deletions apps/server/src/environment/ServerEnvironment.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -244,11 +244,13 @@ it.layer(NodeServices.layer)("ServerEnvironmentLive", (it) => {
expect(withFd.capabilities.serverSelfUpdate).toBe("desktop-managed");
expect(withFd.capabilities.desktopAppUpdate).toBe(true);
expect(withFd.capabilities.serverSelfUpdateProgress).toBe(true);
expect(withFd.capabilities.serverUpdateThreadContinuation).toBe(true);

const withoutFd = yield* describeWith({ mode: "desktop" });
expect(withoutFd.capabilities.serverSelfUpdate).toBe("desktop-managed");
expect(withoutFd.capabilities.desktopAppUpdate).toBeUndefined();
expect(withoutFd.capabilities.serverSelfUpdateProgress).toBeUndefined();
expect(withoutFd.capabilities.serverUpdateThreadContinuation).toBeUndefined();

const web = yield* describeWith({ mode: "web", desktopTelemetryControlFd: 5 });
expect(web.capabilities.desktopAppUpdate).toBeUndefined();
Expand Down
5 changes: 4 additions & 1 deletion apps/server/src/environment/ServerEnvironment.ts
Original file line number Diff line number Diff line change
Expand Up @@ -224,7 +224,10 @@ export const make = Effect.gen(function* () {
threadPullRequestLinking: true,
...(serverSelfUpdate === null ? {} : { serverSelfUpdate }),
...(serverSelfUpdate === "boot-service" || desktopAppUpdate
? { serverSelfUpdateProgress: true }
? {
serverSelfUpdateProgress: true,
serverUpdateThreadContinuation: true,
}
: {}),
...(desktopAppUpdate ? { desktopAppUpdate: true } : {}),
},
Expand Down
1 change: 1 addition & 0 deletions apps/server/src/provider/Layers/CodexAdapter.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2010,6 +2010,7 @@ export const makeCodexAdapter = Effect.fn("makeCodexAdapter")(function* (
provider: PROVIDER,
capabilities: {
sessionModelSwitch: "in-session",
promptlessTurnContinuation: true,
},
startSession,
sendTurn,
Expand Down
51 changes: 51 additions & 0 deletions apps/server/src/provider/Layers/ProviderService.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -246,6 +246,7 @@ function makeFakeCodexAdapter(provider: ProviderDriverKind = CODEX_DRIVER) {
provider,
capabilities: {
sessionModelSwitch: "in-session",
...(provider === CODEX_DRIVER ? { promptlessTurnContinuation: true } : {}),
},
startSession,
sendTurn,
Expand Down Expand Up @@ -960,6 +961,56 @@ it.effect(
);

routing.layer("ProviderServiceLive routing", (it) => {
it.effect("allows promptless continuation only for capable providers", () =>
Effect.gen(function* () {
const provider = yield* ProviderService.ProviderService;
const codexThreadId = asThreadId("thread-promptless-continuation");
yield* provider.startSession(codexThreadId, {
provider: CODEX_DRIVER,
providerInstanceId: codexInstanceId,
threadId: codexThreadId,
runtimeMode: "full-access",
});

yield* provider.sendTurn({ threadId: codexThreadId, continuation: true });
assert.deepEqual(routing.codex.sendTurn.mock.calls.at(-1)?.[0], {
threadId: codexThreadId,
continuation: true,
});

const claudeThreadId = asThreadId("thread-promptless-continuation-unsupported");
yield* provider.startSession(claudeThreadId, {
provider: CLAUDE_AGENT_DRIVER,
providerInstanceId: claudeAgentInstanceId,
threadId: claudeThreadId,
runtimeMode: "full-access",
});
const failure = yield* Effect.flip(
provider.sendTurn({ threadId: claudeThreadId, continuation: true }),
);
assert.instanceOf(failure, ProviderValidationError);
assert.include(failure.issue, "requires an explicit continuation prompt");
assert.equal(routing.claude.sendTurn.mock.calls.length, 0);

yield* provider.stopSession({ threadId: claudeThreadId });
routing.claude.startSession.mockClear();
const stoppedFailure = yield* Effect.flip(
provider.sendTurn({ threadId: claudeThreadId, continuation: true }),
);
assert.instanceOf(stoppedFailure, ProviderValidationError);
assert.include(stoppedFailure.issue, "requires an explicit continuation prompt");
assert.equal(routing.claude.startSession.mock.calls.length, 0);

yield* provider.stopSession({ threadId: codexThreadId });
routing.codex.startSession.mockClear();
routing.codex.sendTurn.mockClear();
routing.codex.stopSession.mockClear();
routing.claude.startSession.mockClear();
routing.claude.sendTurn.mockClear();
routing.claude.stopSession.mockClear();
}),
);

it.effect("routes provider operations and rollback conversation", () =>
Effect.gen(function* () {
const provider = yield* ProviderService.ProviderService;
Expand Down
Loading
Loading