Skip to content
Open
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
31 changes: 30 additions & 1 deletion apps/server/src/mcp/WorktreeMcpService.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,7 @@ import {
} from "@t3tools/contracts";
import * as Cause from "effect/Cause";
import * as Deferred from "effect/Deferred";
import * as Duration from "effect/Duration";
import * as Effect from "effect/Effect";
import * as Exit from "effect/Exit";
import * as Fiber from "effect/Fiber";
Expand Down Expand Up @@ -92,6 +93,7 @@ interface HarnessOptions {
readonly currentBranch?: string | null;
readonly notARepo?: boolean;
readonly newWorktreesStartFromOrigin?: boolean;
readonly automaticGitFetchInterval?: Duration.Duration;
readonly setupScript?: "started" | "no-script" | "fails" | "dies";
readonly dispatchFails?: boolean;
readonly dispatchDies?: boolean;
Expand Down Expand Up @@ -260,7 +262,9 @@ const makeHarness = (options: HarnessOptions = {}) => {
workingTree: { files: [], insertions: 0, deletions: 0 },
}),
);
const refreshStatus = vi.fn((_: string) => Effect.die("refreshStatus stub"));
const refreshStatus = vi.fn<
VcsStatusBroadcaster.VcsStatusBroadcaster["Service"]["refreshStatus"]
>(() => Effect.die("refreshStatus stub"));
const runForThread = vi.fn((input: { readonly worktreePath: string }) => {
switch (options.setupScript ?? "started") {
case "no-script":
Expand Down Expand Up @@ -317,6 +321,9 @@ const makeHarness = (options: HarnessOptions = {}) => {
} satisfies Partial<ProjectService.ProjectService["Service"]>),
ServerSettings.layerTest({
newWorktreesStartFromOrigin: options.newWorktreesStartFromOrigin ?? false,
...(options.automaticGitFetchInterval === undefined
? {}
: { automaticGitFetchInterval: options.automaticGitFetchInterval }),
}),
Layer.mock(GitWorkflowService.GitWorkflowService)({
listRefs,
Expand Down Expand Up @@ -351,6 +358,7 @@ const makeHarness = (options: HarnessOptions = {}) => {
deleteLocalBranch,
localStatus,
runForThread,
refreshStatus,
};
};

Expand Down Expand Up @@ -430,6 +438,27 @@ describe("t3_worktree_handoff", () => {
});
});

it.effect("refreshes the new worktree's status with the configured Git fetch interval", () =>
Effect.gen(function* () {
for (const [configured, expected] of [
[Duration.zero, Duration.zero],
[undefined, Duration.seconds(30)],
] as const) {
const harness = makeHarness(
configured === undefined ? {} : { automaticGitFetchInterval: configured },
);
yield* runHandoff(harness, { branch: "feature/refresh" });

expect(harness.refreshStatus).toHaveBeenCalledTimes(1);
const [cwd, refreshOptions] = harness.refreshStatus.mock.calls[0]!;
expect(cwd).toBe("/worktrees/project/feature/refresh");
const interval = refreshOptions?.automaticRemoteRefreshInterval;
if (interval === undefined) return expect.fail("refreshStatus got no fetch interval");
expect(Duration.toMillis(yield* interval)).toBe(Duration.toMillis(expected));
}
}),
);

it.effect("skips the continuation when no continuationPrompt is given", () => {
const harness = makeHarness();
return Effect.gen(function* () {
Expand Down
15 changes: 14 additions & 1 deletion apps/server/src/mcp/WorktreeMcpService.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
import {
CommandId,
DEFAULT_AUTOMATIC_GIT_FETCH_INTERVAL,
MessageId,
type ProjectId,
WorktreeMcpFailure,
Expand All @@ -9,6 +10,7 @@ import {
type WorktreeMcpSetupScriptStatus,
type WorktreeMcpStatusResult,
} from "@t3tools/contracts";
import { resolveServerBackgroundActivitySettings } from "@t3tools/shared/backgroundActivitySettings";
import * as Cause from "effect/Cause";
import * as Context from "effect/Context";
import * as Crypto from "effect/Crypto";
Expand Down Expand Up @@ -126,6 +128,15 @@ const make = Effect.gen(function* () {
asOperationFailed("Unable to read server settings"),
);

// The post-handoff status refresh honors the Git fetch interval, so a zero
// interval keeps it from fetching, like the client's refresh RPC.
const automaticGitFetchInterval = serverSettings.getSettings.pipe(
Effect.map(
(settings) => resolveServerBackgroundActivitySettings(settings).automaticGitFetchInterval,
),
Effect.orElseSucceed(() => DEFAULT_AUTOMATIC_GIT_FETCH_INTERVAL),
);

const handoffIds = (scope: McpThreadInvocationScope) =>
crypto.randomUUIDv4.pipe(
Effect.map((uuid) => {
Expand Down Expand Up @@ -392,7 +403,9 @@ const make = Effect.gen(function* () {
const continuation = yield* recheckAndBind.pipe(Effect.andThen(queueContinuation));

yield* vcsStatusBroadcaster
.refreshStatus(worktreePath)
.refreshStatus(worktreePath, {
automaticRemoteRefreshInterval: automaticGitFetchInterval,
})
.pipe(Effect.ignoreCause({ log: true }), Effect.forkDetach);

let setupScript: WorktreeMcpSetupScriptStatus = { status: "skipped" };
Expand Down
140 changes: 140 additions & 0 deletions apps/server/src/vcs/VcsStatusBroadcaster.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -212,6 +212,113 @@ describe("VcsStatusBroadcaster", () => {
},
);

it.effect("does not pull automatically when periodic refreshes are disabled", () => {
let remoteStatus: VcsStatusRemoteResult = { ...baseRemoteStatus, behindCount: 2 };
let pullCalls = 0;
const localStatus: VcsStatusLocalResult = {
...baseLocalStatus,
isDefaultRef: true,
refName: "main",
};
const testLayer = VcsStatusBroadcaster.layer.pipe(
Layer.provideMerge(NodeServices.layer),
Layer.provide(layerBackgroundPolicy(() => true)),
Layer.provide(
Layer.succeed(VcsStatusBroadcaster.VcsAutoPullPolicy, {
isEnabled: () => Effect.succeed(true),
}),
),
Layer.provide(
Layer.mock(GitWorkflowService.GitWorkflowService)({
localStatus: () => Effect.succeed(localStatus),
remoteStatus: () => Effect.succeed(remoteStatus),
invalidateLocalStatus: () => Effect.void,
invalidateRemoteStatus: () => Effect.void,
invalidateStatus: () => Effect.void,
pullCurrentBranch: () =>
Effect.sync(() => {
pullCalls += 1;
remoteStatus = { ...remoteStatus, behindCount: 0 };
return { status: "pulled" as const, refName: "main", upstreamRef: "origin/main" };
}),
}),
),
);

return Effect.gen(function* () {
const broadcaster = yield* VcsStatusBroadcaster.VcsStatusBroadcaster;

// The cached upstream still says "behind", but nothing fetched, so a
// focus refresh with the interval at 0 must not reach for the remote.
const quiet = yield* broadcaster.refreshStatus("/repo", {
automaticRemoteRefreshInterval: Effect.succeed(Duration.zero),
});
assert.equal(pullCalls, 0);
assert.equal(quiet.behindCount, 2);

const pulled = yield* broadcaster.refreshStatus("/repo");
assert.equal(pullCalls, 1);
assert.equal(pulled.behindCount, 0);
}).pipe(Effect.provide(testLayer));
});

it.effect("does not pull automatically on the initial poll when the interval is zero", () => {
let pullCalls = 0;
// Settles on the first remote update or pull, whichever the poll reaches.
const settled = Deferred.makeUnsafe<void>();
const behindRemote: VcsStatusRemoteResult = { ...baseRemoteStatus, behindCount: 2 };
const testLayer = VcsStatusBroadcaster.layer.pipe(
Layer.provideMerge(NodeServices.layer),
Layer.provide(layerBackgroundPolicy(() => true)),
Layer.provide(
Layer.succeed(VcsStatusBroadcaster.VcsAutoPullPolicy, {
isEnabled: () => Effect.succeed(true),
}),
),
Layer.provide(
Layer.mock(GitWorkflowService.GitWorkflowService)({
localStatus: () =>
Effect.succeed({ ...baseLocalStatus, isDefaultRef: true, refName: "main" }),
remoteStatus: () => Effect.succeed(behindRemote),
invalidateLocalStatus: () => Effect.void,
invalidateRemoteStatus: () => Effect.void,
invalidateStatus: () => Effect.void,
pullCurrentBranch: () =>
Effect.sync(() => {
pullCalls += 1;
return { status: "pulled" as const, refName: "main", upstreamRef: "origin/main" };
}).pipe(Effect.tap(() => Deferred.succeed(settled, undefined))),
}),
),
);

return Effect.gen(function* () {
const broadcaster = yield* VcsStatusBroadcaster.VcsStatusBroadcaster;
const scope = yield* Scope.make();
let remoteUpdated: VcsStatusStreamEvent | undefined;
yield* Stream.runForEach(
broadcaster.streamStatus(
{ cwd: "/repo" },
{ automaticRemoteRefreshInterval: Effect.succeed(Duration.zero) },
),
(event) => {
if (event._tag !== "remoteUpdated") return Effect.void;
remoteUpdated = event;
return Deferred.succeed(settled, undefined);
},
).pipe(Effect.forkIn(scope));

yield* Deferred.await(settled);
assert.equal(pullCalls, 0);
assert.deepStrictEqual(remoteUpdated, {
_tag: "remoteUpdated",
remote: behindRemote,
} satisfies VcsStatusStreamEvent);

yield* Scope.close(scope, Exit.void);
}).pipe(Effect.provide(testLayer));
});

it.effect("reuses the cached VCS status across repeated reads", () => {
const state = {
currentLocalStatus: baseLocalStatus,
Expand Down Expand Up @@ -594,6 +701,39 @@ describe("VcsStatusBroadcaster", () => {
},
);

it.effect(
"an explicit refresh reads the cached upstream when periodic refreshes are disabled",
() => {
const state = {
currentLocalStatus: baseLocalStatus,
currentRemoteStatus: remoteStatusWithPr,
localStatusCalls: 0,
remoteStatusCalls: 0,
localInvalidationCalls: 0,
remoteInvalidationCalls: 0,
remoteStatusRefreshUpstreamValues: [] as Array<boolean | undefined>,
};

return Effect.gen(function* () {
const broadcaster = yield* VcsStatusBroadcaster.VcsStatusBroadcaster;

// A window focus or a mobile thread selection with the interval at 0.
yield* broadcaster.refreshStatus("/repo", {
automaticRemoteRefreshInterval: Effect.succeed(Duration.zero),
});
assert.deepStrictEqual(state.remoteStatusRefreshUpstreamValues, [false]);
assert.equal(state.remoteInvalidationCalls, 1);

yield* broadcaster.refreshStatus("/repo", {
automaticRemoteRefreshInterval: Effect.succeed(Duration.minutes(1)),
});
yield* broadcaster.refreshStatus("/repo");
assert.deepStrictEqual(state.remoteStatusRefreshUpstreamValues, [false, true, true]);
assert.equal(state.remoteStatusCalls, 3);
}).pipe(Effect.provide(layerTestFor(state)));
},
);

it.effect("passive streams retain cached remote status", () => {
const state = {
currentLocalStatus: baseLocalStatus,
Expand Down
26 changes: 21 additions & 5 deletions apps/server/src/vcs/VcsStatusBroadcaster.ts
Comment thread
Mnigos marked this conversation as resolved.
Original file line number Diff line number Diff line change
Expand Up @@ -192,7 +192,10 @@ export class VcsStatusBroadcaster extends Context.Service<
readonly refreshLocalStatus: (
cwd: string,
) => Effect.Effect<VcsStatusLocalResult, GitManagerServiceError>;
readonly refreshStatus: (cwd: string) => Effect.Effect<VcsStatusResult, GitManagerServiceError>;
readonly refreshStatus: (
cwd: string,
options?: StreamStatusOptions,
) => Effect.Effect<VcsStatusResult, GitManagerServiceError>;
/**
* Refresh a loaded cwd after a turn if background policy allows it.
* GitManager retries missing PRs for the current branch and keeps known
Expand Down Expand Up @@ -459,7 +462,11 @@ export const make = Effect.gen(function* () {
}
const previousRemote = (yield* getCachedStatus(cwd))?.remote?.value;
const remote = yield* workflow.remoteStatus({ cwd }, options);
const pulled = yield* maybeAutoPull(cwd, remote, options?.policyCwds ?? [cwd]);
// Like refreshStatus: no automatic pull when nothing was fetched.
const pulled =
options?.refreshUpstream === false
? null
: yield* maybeAutoPull(cwd, remote, options?.policyCwds ?? [cwd]);
if (pulled !== null) return pulled.remote;
// Local status holds the Changes totals, which compare against remote refs. A fetch can
// move them with no local trigger (a push from a terminal, a PR merged on the host), so
Expand All @@ -480,18 +487,27 @@ export const make = Effect.gen(function* () {

const refreshStatus: VcsStatusBroadcaster["Service"]["refreshStatus"] = Effect.fn(
"VcsStatusBroadcaster.refreshStatus",
)(function* (rawCwd) {
)(function* (rawCwd, options) {
const cwd = yield* withFileSystem(normalizeCwd(rawCwd));
// A zero fetch interval is the user's promise that nothing fetches on its
// own, and a refresh runs on every window focus, so it reads the upstream
// the cache already has instead of fetching it.
const configuredInterval = yield* (
options?.automaticRemoteRefreshInterval ?? Effect.succeed(DEFAULT_VCS_STATUS_REFRESH_INTERVAL)
);
const refreshUpstream = !Duration.isZero(configuredInterval);
// invalidateStatus (not the two partial invalidations) so an explicit
// refresh also bypasses GitManager's slow PR-lookup cache.
return yield* withRemoteWriteLock(
cwd,
Effect.gen(function* () {
yield* workflow.invalidateStatus(cwd);
// Local after remote: the fetch can move the base that the Changes totals compare with.
const remote = yield* workflow.remoteStatus({ cwd });
const remote = yield* workflow.remoteStatus({ cwd }, { refreshUpstream });
const local = yield* workflow.localStatus({ cwd });
const pulled = yield* maybeAutoPull(cwd, remote, [rawCwd]);
// An automatic pull contacts the remote too, and without a fetch the
// "behind" count it would act on is stale anyway.
const pulled = refreshUpstream ? yield* maybeAutoPull(cwd, remote, [rawCwd]) : null;
if (pulled !== null) return mergeGitStatusParts(pulled.local, pulled.remote);
return yield* updateCachedStatus(cwd, local, remote, { publish: true });
}),
Expand Down
7 changes: 5 additions & 2 deletions apps/server/src/ws.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1730,7 +1730,7 @@ const layerWsRpc = (

const refreshGitStatus = (cwd: string) =>
vcsStatusBroadcaster
.refreshStatus(cwd)
.refreshStatus(cwd, { automaticRemoteRefreshInterval: automaticGitFetchInterval })
.pipe(Effect.ignoreCause({ log: true }), Effect.forkDetach, Effect.asVoid);

const getOrchestrationV2ArchivedShellSnapshot = sql
Expand Down Expand Up @@ -2767,7 +2767,10 @@ const layerWsRpc = (
worktreeSetupTracker
.cancel(input.threadId)
.pipe(Effect.map((cancelled) => ({ cancelled }))),
[WS_METHODS.vcsRefreshStatus]: (input) => vcsStatusBroadcaster.refreshStatus(input.cwd),
[WS_METHODS.vcsRefreshStatus]: (input) =>
vcsStatusBroadcaster.refreshStatus(input.cwd, {
automaticRemoteRefreshInterval: automaticGitFetchInterval,
}),
[WS_METHODS.vcsPull]: (input) =>
gitWorkflow.pullCurrentBranch(input.cwd).pipe(
Effect.matchCauseEffect({
Expand Down
Loading