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
12 changes: 12 additions & 0 deletions apps/server/src/provider/Layers/AntigravityProvider.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -620,6 +620,18 @@ it.layer(testLayer)("Antigravity provider snapshots", (it) => {
after.workspaceSnapshots?.find((entry) => entry.cwd === "/workspace")?.skills,
).toEqual(skills);
expect((yield* harness.provider.snapshotForCwd("/workspace")).skills).toEqual(skills);

const rescanned = [
...skills,
{ name: "review", path: "/workspace/.agent/skills/review", enabled: true },
];
yield* harness.provider.snapshotForCwd("/workspace", rescanned);
yield* harness.provider.onSessionStarted(started, "/workspace");
expect(
(yield* harness.provider.snapshot.getSnapshot).workspaceSnapshots?.find(
(entry) => entry.cwd === "/workspace",
)?.skills,
).toEqual(rescanned);
}),
),
);
Expand Down
19 changes: 18 additions & 1 deletion apps/server/src/provider/Layers/AntigravityProvider.ts
Original file line number Diff line number Diff line change
Expand Up @@ -377,7 +377,24 @@ export const makeAntigravityProvider = Effect.fn("makeAntigravityProvider")(func
cwd: string,
skills?: ServerProvider["skills"],
) {
if (skills) discoveredSkills.set(cwd, skills);
if (skills) {
discoveredSkills.set(cwd, skills);
// A rescan replaces the stored entry's skills. Session callbacks and
// health checks republish that entry, so it must not keep old ones.
yield* SubscriptionRef.update(metadata, (state) =>
state.draft.workspaceSnapshots?.some((entry) => entry.cwd === cwd)
? {
...state,
draft: {
...state.draft,
workspaceSnapshots: state.draft.workspaceSnapshots.map((entry) =>
entry.cwd === cwd ? { ...entry, skills } : entry,
),
},
}
: state,
);
}
const snapshot = yield* getSnapshot;
const workspace = snapshot.workspaceSnapshots?.find((entry) => entry.cwd === cwd);
const resolvedSkills = skills ?? workspace?.skills ?? discoveredSkills.get(cwd) ?? [];
Expand Down
58 changes: 57 additions & 1 deletion apps/server/src/provider/Layers/ProviderRegistry.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1479,6 +1479,12 @@ it.layer(Layer.mergeAll(TestNodeServices, ServerSettingsModule.layerTest(), Test
slashCommands: [],
} as const satisfies ServerProvider;
const snapshotCalls = yield* Ref.make(0);
const scopedResult = yield* Ref.make<ServerProvider>(scopedProvider);
const cacheInvalidations = yield* Ref.make(0);
const scanGate = yield* Ref.make<{
readonly started: Deferred.Deferred<void>;
readonly release: Deferred.Deferred<void>;
} | null>(null);
const returnPendingSnapshot = yield* Ref.make(true);
const probeStarted = yield* Deferred.make<void>();
const releaseProbe = yield* Deferred.make<void>();
Expand Down Expand Up @@ -1508,6 +1514,7 @@ it.layer(Layer.mergeAll(TestNodeServices, ServerSettingsModule.layerTest(), Test
applyUsageLimits: () => Effect.void,
},
snapshotForCwd,
invalidateCaches: Ref.update(cacheInvalidations, (count) => count + 1),
adapter: {} as ProviderInstance["adapter"],
textGeneration: {} as ProviderInstance["textGeneration"],
});
Expand All @@ -1517,7 +1524,13 @@ it.layer(Layer.mergeAll(TestNodeServices, ServerSettingsModule.layerTest(), Test
if (yield* Ref.get(returnPendingSnapshot)) return pendingScopedProvider;
yield* Deferred.succeed(probeStarted, undefined);
yield* Deferred.await(releaseProbe);
return scopedProvider;
const result = yield* Ref.get(scopedResult);
const gate = yield* Ref.getAndSet(scanGate, null);
if (gate) {
yield* Deferred.succeed(gate.started, undefined);
yield* Deferred.await(gate.release);
}
return result;
}),
);
const rebuiltProvider = {
Expand Down Expand Up @@ -1593,6 +1606,49 @@ it.layer(Layer.mergeAll(TestNodeServices, ServerSettingsModule.layerTest(), Test
);
yield* registry.refreshWorkspaceSnapshot({ instanceId, cwd: "/workspace" });
assert.strictEqual(yield* Ref.get(snapshotCalls), 2);
const newSkills = [
...scopedProvider.skills,
{ name: "added", path: "/workspace/added/SKILL.md", enabled: true },
];
yield* Ref.set(scopedResult, { ...scopedProvider, skills: newSkills });
yield* registry.refreshWorkspaceSnapshot({
instanceId,
cwd: "/workspace",
fresh: true,
});
assert.strictEqual(yield* Ref.get(snapshotCalls), 3);
assert.strictEqual(yield* Ref.get(cacheInvalidations), 1);
assert.deepStrictEqual(
(yield* registry.getProviders)[0]?.workspaceSnapshots?.map((s) => s.skills),
[newSkills],
);

// A slow fresh scan that read older files must not overwrite a
// newer scan that finished first.
const slowStarted = yield* Deferred.make<void>();
const releaseSlow = yield* Deferred.make<void>();
yield* Ref.set(scanGate, { started: slowStarted, release: releaseSlow });
yield* Ref.set(scopedResult, scopedProvider);
const slowScan = yield* registry
.refreshWorkspaceSnapshot({ instanceId, cwd: "/workspace", fresh: true })
.pipe(Effect.forkChild);
yield* Deferred.await(slowStarted);
const latestSkills = [
...newSkills,
{ name: "latest", path: "/workspace/latest/SKILL.md", enabled: true },
];
yield* Ref.set(scopedResult, { ...scopedProvider, skills: latestSkills });
yield* registry.refreshWorkspaceSnapshot({
instanceId,
cwd: "/workspace",
fresh: true,
});
yield* Deferred.succeed(releaseSlow, undefined);
yield* Fiber.join(slowScan);
assert.deepStrictEqual(
(yield* registry.getProviders)[0]?.workspaceSnapshots?.map((s) => s.skills),
[latestSkills],
);

yield* Ref.set(instancesRef, [rebuiltInstance]);
yield* PubSub.publish(registryChanges, undefined);
Expand Down
93 changes: 67 additions & 26 deletions apps/server/src/provider/Layers/ProviderRegistry.ts
Original file line number Diff line number Diff line change
Expand Up @@ -82,6 +82,15 @@ const hasModelCapabilities = (model: ServerProvider["models"][number]): boolean

const MAX_WORKSPACE_SNAPSHOTS_PER_PROVIDER = 16;

function dropProviderWorkspaceSnapshot(provider: ServerProvider, cwd: string): ServerProvider {
return provider.workspaceSnapshots?.some((snapshot) => snapshot.cwd === cwd)
? {
...provider,
workspaceSnapshots: provider.workspaceSnapshots.filter((snapshot) => snapshot.cwd !== cwd),
}
: provider;
}

export function upsertProviderWorkspaceSnapshot(
provider: ServerProvider,
cwd: string,
Expand Down Expand Up @@ -835,17 +844,44 @@ export const ProviderRegistryLive = Layer.effect(
return yield* Ref.get(providersRef);
});

const updateProviders = (
update: (providers: ReadonlyArray<ServerProvider>) => ReadonlyArray<ServerProvider>,
) =>
Ref.modify(providersRef, (currentProviders) => {
const nextProviders = update(currentProviders);
return [[currentProviders, nextProviders] as const, nextProviders];
}).pipe(
Effect.tap(([previousProviders, nextProviders]) =>
haveProvidersChanged(previousProviders, nextProviders)
? PubSub.publish(changesPubSub, nextProviders)
: Effect.void,
),
Effect.map(([, nextProviders]) => nextProviders),
);

const refreshWorkspaceSnapshot = Effect.fn("refreshWorkspaceSnapshot")(function* (input: {
readonly instanceId: ProviderInstanceId;
readonly cwd: string;
readonly fresh?: boolean;
}) {
// Fresh scans drop other instances' snapshots for this cwd first, so a
// composer on one of them scans again on next use, even when this
// instance is gone or cannot be scanned.
if (input.fresh) {
yield* updateProviders((providers) =>
providers.map((candidate) =>
candidate.instanceId === input.instanceId
? candidate
: dropProviderWorkspaceSnapshot(candidate, input.cwd),
),
);
}
const providers = yield* Ref.get(providersRef);
const provider = providers.find((candidate) => candidate.instanceId === input.instanceId);
if (
!provider ||
!provider.enabled ||
provider.workspaceSnapshots?.some((s) => s.cwd === input.cwd)
) {
const workspaceSnapshotOf = (candidate: ServerProvider | undefined) =>
candidate?.workspaceSnapshots?.find((s) => s.cwd === input.cwd);
const scannedFrom = workspaceSnapshotOf(provider);
if (!provider || !provider.enabled || (!input.fresh && scannedFrom)) {
return providers;
}
const instance = yield* instanceRegistry.getInstance(input.instanceId);
Expand All @@ -857,42 +893,47 @@ export const ProviderRegistryLive = Layer.effect(
next.set(instance, new Set(current).add(input.cwd));
return [true, next] as const;
});
if (!claimed) return yield* Ref.get(providersRef);
return yield* instance.snapshotForCwd(input.cwd).pipe(
// A fresh scan never joins a running one, which may predate the change.
if (!claimed && !input.fresh) return yield* Ref.get(providersRef);
// Fresh scans also re-read the machine snapshot: Claude's plugin
// commands come from it, not from the cwd scan.
const refreshMachineSnapshot = input.fresh
? (instance.invalidateCaches ?? Effect.void).pipe(
Comment thread
t3dotgg marked this conversation as resolved.
Effect.andThen(refreshInstance(input.instanceId)),
)
: Effect.void;
return yield* refreshMachineSnapshot.pipe(
Effect.andThen(instance.snapshotForCwd(input.cwd)),
Effect.flatMap((scopedSnapshot) =>
scopedSnapshot.status === "error"
? Ref.get(providersRef)
: instanceRegistry.getInstance(input.instanceId).pipe(
Effect.flatMap((currentInstance) => {
if (currentInstance !== instance) return Ref.get(providersRef);
return Ref.modify(providersRef, (currentProviders) => {
const nextProviders = currentProviders.map((candidate) =>
// Write only if the cwd's snapshot did not change during the
// scan. A session event or another scan that landed first is newer.
return updateProviders((currentProviders) =>
currentProviders.map((candidate) =>
candidate.instanceId === input.instanceId &&
!candidate.workspaceSnapshots?.some((s) => s.cwd === input.cwd)
Equal.equals(workspaceSnapshotOf(candidate), scannedFrom)
? upsertProviderWorkspaceSnapshot(candidate, input.cwd, scopedSnapshot)
: candidate,
);
return [[currentProviders, nextProviders] as const, nextProviders];
}).pipe(
Effect.tap(([previousProviders, nextProviders]) =>
haveProvidersChanged(previousProviders, nextProviders)
? PubSub.publish(changesPubSub, nextProviders)
: Effect.void,
),
Effect.map(([, nextProviders]) => nextProviders),
);
}),
),
),
Effect.ensuring(
Ref.update(workspaceRefreshesRef, (refreshes) => {
const next = new Map(refreshes);
const current = new Set(next.get(instance));
current.delete(input.cwd);
if (current.size) next.set(instance, current);
else next.delete(instance);
return next;
}),
claimed
? Ref.update(workspaceRefreshesRef, (refreshes) => {
const next = new Map(refreshes);
const current = new Set(next.get(instance));
current.delete(input.cwd);
if (current.size) next.set(instance, current);
else next.delete(instance);
return next;
})
: Effect.void,
),
);
});
Expand Down
7 changes: 7 additions & 0 deletions apps/server/src/provider/Services/ProviderRegistry.ts
Original file line number Diff line number Diff line change
Expand Up @@ -48,9 +48,16 @@ export interface ProviderRegistryShape {
instanceId: ProviderInstanceId,
) => Effect.Effect<ReadonlyArray<ServerProvider>>;

/**
* Fill the skills and slash commands snapshot for one cwd. A cwd that
* already has a snapshot is left alone unless `fresh` is set. A fresh scan
* also refreshes the instance's machine snapshot and drops other
* instances' snapshots for the cwd.
*/
readonly refreshWorkspaceSnapshot: (input: {
readonly instanceId: ProviderInstanceId;
readonly cwd: string;
readonly fresh?: boolean;
}) => Effect.Effect<ReadonlyArray<ServerProvider>>;

/**
Expand Down
1 change: 1 addition & 0 deletions apps/server/src/ws.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2612,6 +2612,7 @@ const makeWsRpcLayer = (
? providerRegistry.refreshWorkspaceSnapshot({
instanceId: input.instanceId,
cwd: input.cwd,
fresh: input.fresh === true,
})
: input.instanceId !== undefined
? providerRegistry.refreshInstance(input.instanceId)
Expand Down
52 changes: 52 additions & 0 deletions apps/web/src/components/CommandPalette.tsx
Original file line number Diff line number Diff line change
Expand Up @@ -53,6 +53,7 @@ import {
MonitorIcon,
MoonIcon,
PaletteIcon,
RotateCcwIcon,
SettingsIcon,
SquarePenIcon,
SunIcon,
Expand Down Expand Up @@ -93,6 +94,8 @@ import { desktopLocalBackendId } from "../connection/desktopLocal";
import { filesystemEnvironment } from "../state/filesystem";
import { projectEnvironment } from "../state/projects";
import { useEnvironmentQuery } from "../state/query";
import { serverEnvironment } from "../state/server";
import { threadEnvironment } from "../state/threads";
import { sourceControlEnvironment } from "../state/sourceControl";
import { useAtomCommand } from "../state/use-atom-command";
import { useAtomQueryRunner } from "../state/use-atom-query-runner";
Expand Down Expand Up @@ -725,6 +728,12 @@ function OpenCommandPaletteDialog(props: {
const startProjectClone = useAtomCommand(sourceControlEnvironment.startProjectClone, {
reportFailure: false,
});
const stopThreadSession = useAtomCommand(threadEnvironment.stopSession, {
reportFailure: false,
});
const refreshProviders = useAtomCommand(serverEnvironment.refreshProviders, {
reportFailure: false,
});
const { environments } = useEnvironments();
const desktopLocalBootstraps = useDesktopLocalBootstraps();
const primaryEnvironmentId = usePrimaryEnvironmentId();
Expand Down Expand Up @@ -1871,6 +1880,49 @@ function OpenCommandPaletteDialog(props: {
}
}

if (activeThread !== null) {
const thread = activeThread;
actionItems.push({
kind: "action",
value: "action:restart-agent-session",
searchTerms: ["restart", "reset", "reload", "agent", "session", "skills", "plugins", "mcp"],
title: "Restart agent session",
icon: <RotateCcwIcon className={ITEM_ICON_CLASS} />,
// Stopping the provider process keeps the conversation: the next message
// spawns a fresh one that resumes it and reloads skills, plugins, and MCP
// servers. The fresh workspace scan updates the composer's slash menu.
// Failures throw into executeItem's error toast.
run: async () => {
const { environmentId } = thread;
if (thread.session && thread.session.status !== "stopped") {
const stopped = await stopThreadSession({
environmentId,
input: { threadId: thread.id },
});
if (stopped._tag === "Failure") throw squashAtomCommandFailure(stopped);
}
// The server stops the process after accepting the command. A failed
// stop shows in the thread.
toastManager.add({
type: "success",
title: "Agent session will restart",
description: "Your next message starts a fresh session.",
});
const project = projectByKey.get(`${environmentId}:${thread.projectId}`);
if (!project) return;
const refreshed = await refreshProviders({
environmentId,
input: {
instanceId: thread.session?.providerInstanceId ?? thread.modelSelection.instanceId,
cwd: thread.worktreePath ?? project.workspaceRoot,
fresh: true,
},
});
if (refreshed._tag === "Failure") throw squashAtomCommandFailure(refreshed);
},
});
}

actionItems.push({
kind: "action",
value: "action:open-file-picker",
Expand Down
4 changes: 4 additions & 0 deletions docs/user/composer.md
Original file line number Diff line number Diff line change
Expand Up @@ -152,6 +152,10 @@ provider. On mobile, both are also available before starting a thread on
The slash menu also includes skills unless you turn off **Settings → General →
Show skills in slash menu**. Only skills enabled for the provider are listed.

After you add or change skills, plugins, or MCP servers, use **Restart agent
session** in the command palette on web and desktop. The conversation continues,
and your next message starts the agent again with the new setup.

Provider commands must start the message to run. T3 Code commands such as
`/model` and `/plan`, and skill mentions, work on any line.

Expand Down
1 change: 1 addition & 0 deletions packages/client-runtime/src/state/server.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1103,6 +1103,7 @@ export function createServerEnvironmentAtoms<R, E>(
environmentId,
input.instanceId ?? null,
input.cwd ?? null,
input.fresh ?? false,
input.refreshModels ?? false,
]),
},
Expand Down
3 changes: 3 additions & 0 deletions packages/contracts/src/rpc.ts
Original file line number Diff line number Diff line change
Expand Up @@ -494,6 +494,9 @@ const WsServerRefreshProvidersRpc = Rpc.make(WS_METHODS.serverRefreshProviders,
*/
instanceId: Schema.optional(ProviderInstanceId),
cwd: Schema.optional(TrimmedNonEmptyString),
/** With `instanceId` and `cwd`: rescan the workspace's skills and slash
* commands even when a snapshot for that cwd already exists. */
fresh: Schema.optional(Schema.Boolean),
/** Explicit user request: bypass T3-owned caches and rediscover models.
* Background status refreshes must not open agent sessions. */
refreshModels: Schema.optional(Schema.Boolean),
Expand Down
Loading