Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
19 commits
Select commit Hold shift + click to select a range
1fe8efd
refactor(web): register the Diff side panel
saphid Oct 2, 2026
bcbd324
refactor(web): register the Preview side panel
github-actions[bot] Oct 4, 2026
6befef3
perf(web): keep the panel host stable across chat view renders
github-actions[bot] Oct 7, 2026
1b23a3e
fix(web): send a late annotation only to the thread it was picked in
github-actions[bot] Oct 7, 2026
16877e3
fix(web): lend the annotation sender only after the chat view commits
github-actions[bot] Oct 7, 2026
bb246d6
refactor(web): move the thread terminal panel and drawer out of ChatView
github-actions[bot] Oct 4, 2026
64cdb14
refactor(web): open the right-panel terminal through the panel host
github-actions[bot] Oct 4, 2026
22f34a8
fix(web): keep a local-checkout terminal launch off the thread's work…
github-actions[bot] Oct 6, 2026
e34c8a4
fix(web): keep a terminal the server opened on the checkout off the t…
github-actions[bot] Oct 7, 2026
f8da909
fix(web): keep device results in the thread that started them
github-actions[bot] Oct 4, 2026
aea3b7a
refactor(web): open the Device side panel through the panel host
github-actions[bot] Oct 4, 2026
df78494
refactor(web): open the pull request side panels through the panel host
github-actions[bot] Oct 4, 2026
d1554f5
refactor(web): open the Files side panel through the panel host
github-actions[bot] Oct 4, 2026
7f84c56
fix(web): drop Files actions that settle after leaving the thread
github-actions[bot] Oct 4, 2026
8136ad4
fix(web): open tree clicks in the thread the Files panel shows now
github-actions[bot] Oct 6, 2026
6d65c39
fix(web): keep a Files browser open across composer draft switches
github-actions[bot] Oct 6, 2026
3be1e97
fix(web): point tree clicks at the new thread as soon as it commits
github-actions[bot] Oct 10, 2026
3d2de4d
feat(pi): show extension statuses in the thread header
github-actions[bot] Oct 4, 2026
890b0c5
fix(pi): keep statuses when Pi refuses a rollback fork
github-actions[bot] Oct 6, 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
1 change: 1 addition & 0 deletions apps/server/src/auth/RpcAuthorization.ts
Original file line number Diff line number Diff line change
Expand Up @@ -167,6 +167,7 @@ export const RPC_REQUIRED_SCOPES = {
[WS_METHODS.subscribeWorktreeSetup]: AuthOrchestrationReadScope,
[WS_METHODS.worktreeSetupCancel]: AuthOrchestrationOperateScope,
[WS_METHODS.subscribeResourceTelemetry]: AuthDiagnosticsReadScope,
[WS_METHODS.subscribeContributionStatus]: AuthOrchestrationReadScope,
[WS_METHODS.vcsRefreshStatus]: AuthOrchestrationReadScope,
[WS_METHODS.gitResolvePullRequest]: AuthOrchestrationReadScope,
[WS_METHODS.vcsListRefs]: AuthOrchestrationReadScope,
Expand Down
206 changes: 206 additions & 0 deletions apps/server/src/contributions/ContributionStatusRpc.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,206 @@
import * as NodeServices from "@effect/platform-node/NodeServices";
import { assert, describe, it } from "@effect/vitest";
import {
AuthOrchestrationReadScope,
AuthRelayReadScope,
type AuthEnvironmentScope,
ProviderDriverKind,
ProviderInstanceId,
ProviderSessionId,
ThreadId,
WS_METHODS,
WsRpcGroup,
} from "@t3tools/contracts";
import * as Deferred from "effect/Deferred";
import * as Effect from "effect/Effect";
import * as Layer from "effect/Layer";
import * as Queue from "effect/Queue";
import * as TestClock from "effect/testing/TestClock";
import { RpcMessage, RpcSerialization, RpcServer } from "effect/rpc";

import * as RpcAuthorization from "../auth/RpcAuthorization.ts";
import { RPC_REQUIRED_SCOPES } from "../auth/RpcAuthorization.ts";
import * as ProviderEventLoggers from "@t3tools/provider-core/server/ProviderEventLoggers";
import * as ProviderOrchestrationAdapterInfrastructure from "../provider/ProviderOrchestrationAdapterInfrastructure.ts";
import * as ContributionStatusStore from "@t3tools/provider-core/server/ContributionStatusStore";

const TAG = WS_METHODS.subscribeContributionStatus;
const THREAD = ThreadId.make("thread-1");
const SOURCE = {
kind: "provider-session",
providerSessionId: ProviderSessionId.make("session-1"),
providerInstanceId: ProviderInstanceId.make("pi"),
driver: ProviderDriverKind.make("pi"),
} as const;
const encodedEntry = (texts: ReadonlyArray<string>) => ({
threadId: THREAD,
source: SOURCE,
items: texts.map((text, index) => ({ key: `k${index}`, text })),
});

// The server group narrowed to this RPC; it keeps the group's scope middleware.
const group = WsRpcGroup.omit(
...[...WsRpcGroup.requests.keys()].filter(
(tag): tag is Exclude<keyof typeof RPC_REQUIRED_SCOPES, typeof TAG> => tag !== TAG,
),
);

/**
* Serves `subscribeContributionStatus` as ws.ts does (the store's subscription
* stream behind the connection's scope middleware) over an in-memory RPC protocol.
*/
const serveStatus = Effect.fn("serveStatus")(function* (
store: ContributionStatusStore.ContributionStatusStoreShape,
scopes: ReadonlyArray<AuthEnvironmentScope>,
) {
const responses = yield* Queue.unbounded<RpcMessage.FromServerEncoded>();
const receive = yield* Deferred.make<Parameters<RpcServer.Protocol["Service"]["run"]>[0]>();
const protocol = yield* RpcServer.Protocol.make((write) =>
Effect.gen(function* () {
yield* Deferred.succeed(receive, write);
const serialization = yield* RpcSerialization.RpcSerialization;
return {
disconnects: yield* Queue.unbounded<number>(),
send: (_clientId, response) => Queue.offer(responses, response),
end: () => Effect.void,
clientIds: Effect.succeed(new Set([0])),
initialMessage: Effect.succeedNone,
supportsAck: false,
supportsTransferables: false,
supportsSpanPropagation: false,
supportsNotifications: true,
codecFor: serialization.codecFor,
};
}),
);
yield* RpcServer.make(group).pipe(
Effect.provide(
Layer.merge(
group.toLayerHandler(TAG, () => ContributionStatusStore.subscriptionStream(store)),
RpcAuthorization.layer(scopes),
),
),
Effect.provideService(RpcServer.Protocol, protocol),
Effect.forkScoped,
);
const write = yield* Deferred.await(receive);
return {
subscribe: (id: string) =>
write(0, { _tag: "Request", id, tag: TAG, payload: {}, headers: [] }),
interrupt: (id: string) => write(0, { _tag: "Interrupt", requestId: id }),
/** The next frame the client receives; the store publishes one per visible change. */
next: Queue.take(responses),
/** Lets server fibers run so a frame that would be sent is in the queue. */
pending: TestClock.adjust(0).pipe(Effect.andThen(Queue.size(responses))),
};
});

const chunk = (
requestId: string,
entries: ReadonlyArray<ReturnType<typeof encodedEntry>>,
): RpcMessage.FromServerEncoded => ({ _tag: "Chunk", requestId, values: [{ entries }] });

describe("subscribeContributionStatus", () => {
it.effect("streams the current snapshot, replacements, and a fresh snapshot on resubscribe", () =>
Effect.gen(function* () {
const store = yield* ContributionStatusStore.make();
const handle = yield* store.openSource(SOURCE);
yield* handle.bindThread(THREAD);
yield* handle.set({ key: "k0", text: "plan" });
const client = yield* serveStatus(store, [AuthOrchestrationReadScope]);

yield* client.subscribe("1");
assert.deepStrictEqual(yield* client.next, chunk("1", [encodedEntry(["plan"])]));

yield* handle.set({ key: "k0", text: "build" });
assert.deepStrictEqual(yield* client.next, chunk("1", [encodedEntry(["build"])]));
yield* handle.clear("k0");
assert.deepStrictEqual(yield* client.next, chunk("1", []));

// A reconnect is a new subscription whose first frame is the current state.
yield* handle.set({ key: "k0", text: "review" });
assert.deepStrictEqual(yield* client.next, chunk("1", [encodedEntry(["review"])]));
yield* client.subscribe("2");
assert.deepStrictEqual(yield* client.next, chunk("2", [encodedEntry(["review"])]));

// An interrupted subscription stops receiving frames; the other keeps going.
yield* client.interrupt("1");
assert.strictEqual((yield* client.next)._tag, "Exit");
yield* handle.set({ key: "k0", text: "done" });
assert.deepStrictEqual(yield* client.next, chunk("2", [encodedEntry(["done"])]));
assert.strictEqual(yield* client.pending, 0);
}).pipe(Effect.provide(RpcSerialization.layerJson), Effect.scoped),
);

it.effect("refuses a client without the orchestration read scope", () =>
Effect.gen(function* () {
const store = yield* ContributionStatusStore.make();
const handle = yield* store.openSource(SOURCE);
yield* handle.bindThread(THREAD);
yield* handle.set({ key: "k0", text: "plan" });
const client = yield* serveStatus(store, [AuthRelayReadScope]);

yield* client.subscribe("1");
assert.deepStrictEqual(yield* client.next, {
_tag: "Exit",
requestId: "1",
exit: {
_tag: "Failure",
cause: [
{
_tag: "Fail",
error: {
_tag: "EnvironmentAuthorizationError",
message: `The authenticated token is missing required scope: ${AuthOrchestrationReadScope}.`,
requiredPermission: AuthOrchestrationReadScope,
requiredScope: AuthOrchestrationReadScope,
},
},
],
},
});
}).pipe(Effect.provide(RpcSerialization.layerJson), Effect.scoped),
);

it.effect("shares one store between provider adapters and the WebSocket stream", () =>
Effect.gen(function* () {
// Mirrors server.ts: adapters get the store through the provider
// infrastructure inside an unwrapped instance-registry layer, while the
// WebSocket layer reads it from the runtime's own reference.
const producer = Layer.unwrap(
Effect.succeed(
Layer.effectDiscard(
Effect.gen(function* () {
const store = yield* ContributionStatusStore.ContributionStatusStore;
const handle = yield* store.openSource(SOURCE);
yield* handle.bindThread(THREAD);
yield* handle.set({ key: "k0", text: "from adapter" });
}),
).pipe(Layer.provide(ProviderOrchestrationAdapterInfrastructure.layer)),
),
);
// As in server.ts, the registry layer sits below the store reference, so
// the producer only sees the store its own infrastructure provides.
const runtime = ContributionStatusStore.layer.pipe(
Layer.provideMerge(producer),
Layer.provide(
Layer.merge(
NodeServices.layer,
Layer.succeed(
ProviderEventLoggers.ProviderEventLoggers,
ProviderEventLoggers.NoOpProviderEventLoggers,
),
),
),
);

const snapshot = yield* Effect.gen(function* () {
const store = yield* ContributionStatusStore.ContributionStatusStore;
return yield* store.snapshot;
}).pipe(Effect.provide(runtime));
assert.deepStrictEqual(snapshot.entries, [
{ threadId: THREAD, source: SOURCE, items: [{ key: "k0", text: "from adapter" }] },
]);
}),
);
});
1 change: 1 addition & 0 deletions apps/server/src/environment/ServerEnvironment.ts
Original file line number Diff line number Diff line change
Expand Up @@ -251,6 +251,7 @@ export const make = Effect.gen(function* () {
serverResolvedCommandContext: true,
environmentIcon: true,
projectCloneTracking: true,
contributionStatus: true,
...(serverSelfUpdate === null ? {} : { serverSelfUpdate }),
...(serverInstallation === null ? {} : { serverInstallation }),
// V2 restart recovery uses the environment-owned opt-in. The old
Expand Down
1 change: 1 addition & 0 deletions apps/server/src/observability/RpcInstrumentation.ts
Original file line number Diff line number Diff line change
Expand Up @@ -154,6 +154,7 @@ const RPC_AGGREGATES = {
[WS_METHODS.subscribeWorktreeSetup]: "vcs",
[WS_METHODS.worktreeSetupCancel]: "vcs",
[WS_METHODS.subscribeResourceTelemetry]: "server",
[WS_METHODS.subscribeContributionStatus]: "server",
[WS_METHODS.vcsRefreshStatus]: "vcs",
[WS_METHODS.vcsPull]: "git",
[WS_METHODS.gitRunStackedAction]: "vcs",
Expand Down
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
import * as Layer from "effect/Layer";

import * as ContributionStatusStore from "@t3tools/provider-core/server/ContributionStatusStore";
import * as ClaudeAdapterV2 from "../orchestration-v2/Adapters/ClaudeAdapterV2.ts";
import * as CodexAdapterV2 from "../orchestration-v2/Adapters/CodexAdapterV2.ts";
import * as CursorAgentSdk from "@t3tools/provider-cursor/server/CursorAgentSdk";
Expand All @@ -20,7 +21,8 @@ export type ProviderOrchestrationAdapterInfrastructure =
* Infrastructure shared by the V2 adapters materialized inside provider
* instances. `providerContinuationRequestsLayer` must be the same layer
* reference the orchestration runtime provides to its continuation worker so
* Effect layer memoization yields one shared queue.
* Effect layer memoization yields one shared queue. The contribution status
* store follows the same rule with the WebSocket server that streams it.
*/
export const layer = Layer.mergeAll(
ClaudeAdapterV2.layerQueryRunner,
Expand All @@ -29,4 +31,5 @@ export const layer = Layer.mergeAll(
CursorKeychain.layer,
IdAllocator.layer,
ProviderContinuationRequests.layer,
ContributionStatusStore.layer,
);
4 changes: 4 additions & 0 deletions apps/server/src/server.ts
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,7 @@ import * as HttpApiBuilder from "effect/http-api/HttpApiBuilder";
import * as BackgroundPolicy from "./background/BackgroundPolicy.ts";
import * as HostPowerMonitor from "./background/HostPowerMonitor.ts";
import * as ServerConfig from "./config.ts";
import * as ContributionStatusStore from "@t3tools/provider-core/server/ContributionStatusStore";
import { withUntracedRequests } from "./http.ts";
import * as ServerHttp from "./http.ts";
import { guardHttpResponseWriteErrors } from "./httpResponseErrorGuard.ts";
Expand Down Expand Up @@ -587,6 +588,9 @@ const layerRuntimeCoreDependenciesBase = Layer.mergeAll(
Layer.provideMerge(layerGit),
Layer.provideMerge(layerVcs),
Layer.provideMerge(Layer.mergeAll(layerTerminal, layerPreview, layerDevice)),
// The same layer reference provider adapters write through, so memoization
// gives producers and the WebSocket stream one store.
Layer.provideMerge(ContributionStatusStore.layer),
Layer.provideMerge(layerPersistence),
// Both read a user-owned file out of the state directory and stream changes
// to clients; neither depends on the other.
Expand Down
6 changes: 5 additions & 1 deletion apps/server/src/ws.ts
Original file line number Diff line number Diff line change
Expand Up @@ -89,7 +89,7 @@ import {
ChatAttachmentId,
PersistChatAttachmentsError,
RpcClientId,
EnvironmentAuthorizationError,
type EnvironmentAuthorizationError,
type ProjectId,
type ProviderDriverKind,
type ProviderInstanceId,
Expand Down Expand Up @@ -212,6 +212,7 @@ import * as DirectEndpoints from "./environment/DirectEndpoints.ts";
import * as RemoteOpenTargets from "./environment/RemoteOpenTargets.ts";
import * as DefectReporter from "./observability/DefectReporter.ts";
import * as BackgroundPolicy from "./background/BackgroundPolicy.ts";
import * as ContributionStatusStore from "@t3tools/provider-core/server/ContributionStatusStore";
import * as EnvironmentAuth from "./auth/EnvironmentAuth.ts";
import { requiredScopeForDeviceList, rpcAuthorizationError } from "./auth/RpcAuthorization.ts";
import * as RpcAuthorization from "./auth/RpcAuthorization.ts";
Expand Down Expand Up @@ -1322,6 +1323,7 @@ const layerWsRpc = (
const hostResources = yield* HostResources.HostResources;
const processResourceMonitor = yield* ProcessResourceMonitor.ProcessResourceMonitor;
const resourceTelemetry = yield* ResourceTelemetry.ResourceTelemetry;
const contributionStatus = yield* ContributionStatusStore.ContributionStatusStore;
const relayClient = yield* RelayClient.RelayClient;
// A webhook URL starts agent runs, so only sessions that may operate
// see it; read-only sessions still see the task itself.
Expand Down Expand Up @@ -3113,6 +3115,8 @@ const layerWsRpc = (
Stream.concat(Stream.make(latest), changes),
),
),
[WS_METHODS.subscribeContributionStatus]: (_input) =>
ContributionStatusStore.subscriptionStream(contributionStatus),
});
return handlers;
}),
Expand Down
15 changes: 15 additions & 0 deletions apps/web/src/browser/openFileInPreview.ts
Original file line number Diff line number Diff line change
Expand Up @@ -54,6 +54,16 @@ export type OpenPreviewMutation<E = unknown> = (input: {
readonly input: PreviewOpenInput;
}) => Promise<AtomCommandResult<PreviewSessionSnapshot, E>>;

/**
* False once the caller has left the thread it started in. Work that has not
* asked the server for a browser by then ends as an interruption. A browser the
* server already opened is still applied to the thread it was opened for, so
* its session never lingers unseen.
*/
type ScopeCheck = (() => boolean) | undefined;

const leftScope = (isScopeCurrent: ScopeCheck) => isScopeCurrent?.() === false;

export async function openUrlInPreview<E>(input: {
readonly threadRef: ScopedThreadRef;
readonly url: string;
Expand All @@ -62,10 +72,12 @@ export async function openUrlInPreview<E>(input: {
readonly profileId?: PreviewOpenInput["profileId"];
/** Open the tab without switching the thread to it. */
readonly background?: boolean;
readonly isScopeCurrent?: ScopeCheck;
}): Promise<AtomCommandResult<void, E | BrowserSettingsReadError>> {
const defaults = await resolveBrowserDefaults().catch(
(cause: unknown) => new BrowserSettingsReadError({ cause }),
);
if (leftScope(input.isScopeCurrent)) return AsyncResult.failure(Cause.interrupt());
if (defaults instanceof BrowserSettingsReadError) {
return AsyncResult.failure(Cause.fail(defaults));
}
Expand Down Expand Up @@ -124,6 +136,7 @@ export async function openFileInPreview<AssetError, PreviewError>(input: {
readonly input: { readonly resource: AssetResource };
}) => Promise<AtomCommandResult<AssetCreateUrlResult, AssetError>>;
readonly openPreview: OpenPreviewMutation<PreviewError>;
readonly isScopeCurrent?: ScopeCheck;
}): Promise<
AtomCommandResult<
void,
Expand Down Expand Up @@ -151,6 +164,7 @@ export async function openFileInPreview<AssetError, PreviewError>(input: {
},
},
});
if (leftScope(input.isScopeCurrent)) return AsyncResult.failure(Cause.interrupt());
if (assetResult._tag === "Failure") {
return AsyncResult.failure(assetResult.cause);
}
Expand All @@ -164,5 +178,6 @@ export async function openFileInPreview<AssetError, PreviewError>(input: {
threadRef: input.threadRef,
url: assetUrl,
openPreview: input.openPreview,
isScopeCurrent: input.isScopeCurrent,
});
}
Loading
Loading