diff --git a/apps/server/src/auth/RpcAuthorization.ts b/apps/server/src/auth/RpcAuthorization.ts index ca55e4e95e01..36371713a7a8 100644 --- a/apps/server/src/auth/RpcAuthorization.ts +++ b/apps/server/src/auth/RpcAuthorization.ts @@ -97,6 +97,10 @@ export const RPC_REQUIRED_SCOPES = { [WS_METHODS.sourceControlLookupRepository]: AuthOrchestrationReadScope, [WS_METHODS.sourceControlCloneRepository]: AuthOrchestrationOperateScope, [WS_METHODS.sourceControlPublishRepository]: AuthOrchestrationOperateScope, + [WS_METHODS.projectCloneStart]: AuthOrchestrationOperateScope, + [WS_METHODS.projectCloneCancel]: AuthOrchestrationOperateScope, + [WS_METHODS.projectCloneRetry]: AuthOrchestrationOperateScope, + [WS_METHODS.subscribeProjectClones]: AuthOrchestrationReadScope, [WS_METHODS.projectsListEntries]: AuthOrchestrationReadScope, [WS_METHODS.projectsReadFile]: AuthOrchestrationReadScope, [WS_METHODS.projectsSearchContents]: AuthOrchestrationReadScope, diff --git a/apps/server/src/bin.test.ts b/apps/server/src/bin.test.ts index 5ba2f281126c..c1bea133f718 100644 --- a/apps/server/src/bin.test.ts +++ b/apps/server/src/bin.test.ts @@ -39,6 +39,7 @@ import * as ProjectionSnapshotQuery from "./orchestration/Services/ProjectionSna import * as OrchestrationEngine from "./orchestration/Services/OrchestrationEngine.ts"; import { OrchestrationLayerLive } from "./orchestration/runtimeLayer.ts"; import { orchestrationHttpApiLayer } from "./orchestration/http.ts"; +import * as ProjectCloneTracker from "./project/ProjectCloneTracker.ts"; import { layerConfig as SqlitePersistenceLayerLive } from "./persistence/Layers/Sqlite.ts"; import * as RepositoryIdentityResolver from "./project/RepositoryIdentityResolver.ts"; import { @@ -363,7 +364,16 @@ const withLiveProjectCliServer = (baseDir: string, run: () => Effect.Ef Effect.gen(function* () { const config = yield* makeCliTestServerConfig(baseDir); const routesLayer = HttpApiBuilder.layer(ProjectCliHttpApi).pipe( - Layer.provide(orchestrationHttpApiLayer), + Layer.provide( + orchestrationHttpApiLayer.pipe( + Layer.provide( + Layer.mock(ProjectCloneTracker.ProjectCloneTracker)({ + get: () => Effect.succeed(null), + discard: () => Effect.void, + }), + ), + ), + ), Layer.provide(environmentAuthenticatedAuthLayer), ); const appLayer = HttpRouter.serve(routesLayer, { diff --git a/apps/server/src/environment/ServerEnvironment.ts b/apps/server/src/environment/ServerEnvironment.ts index 92e32d005492..9c25767245df 100644 --- a/apps/server/src/environment/ServerEnvironment.ts +++ b/apps/server/src/environment/ServerEnvironment.ts @@ -236,6 +236,7 @@ export const make = Effect.gen(function* () { pullRequestStackActions: true, threadPullRequestLinking: true, environmentIcon: true, + projectCloneTracking: true, ...(serverSelfUpdate === null ? {} : { serverSelfUpdate }), ...(serverSelfUpdate === "boot-service" || desktopAppUpdate ? { diff --git a/apps/server/src/orchestration/http.ts b/apps/server/src/orchestration/http.ts index f7147106c7a9..77a04442ccfb 100644 --- a/apps/server/src/orchestration/http.ts +++ b/apps/server/src/orchestration/http.ts @@ -16,6 +16,7 @@ import { failEnvironmentNotFound, requireEnvironmentScope, } from "../auth/http.ts"; +import * as ProjectCloneTracker from "../project/ProjectCloneTracker.ts"; import { OrchestrationEngineService } from "./Services/OrchestrationEngine.ts"; import { ProjectionSnapshotQuery } from "./Services/ProjectionSnapshotQuery.ts"; @@ -25,6 +26,7 @@ export const orchestrationHttpApiLayer = HttpApiBuilder.group( Effect.fnUntraced(function* (handlers) { const projectionSnapshotQuery = yield* ProjectionSnapshotQuery; const orchestrationEngine = yield* OrchestrationEngineService; + const projectCloneTracker = yield* ProjectCloneTracker.ProjectCloneTracker; return handlers .handle( @@ -93,10 +95,18 @@ export const orchestrationHttpApiLayer = HttpApiBuilder.group( Effect.fn("environment.orchestration.dispatch")(function* (args) { yield* annotateEnvironmentRequest(args.endpoint.name); yield* requireEnvironmentScope(AuthOrchestrationOperateScope); + yield* ProjectCloneTracker.rejectCommandsDuringClone( + projectCloneTracker, + args.payload, + ).pipe( + Effect.catch((cause) => + failEnvironmentInternal("orchestration_dispatch_failed", cause), + ), + ); const normalizedCommand = yield* normalizeDispatchCommand(args.payload).pipe( Effect.catch(() => failEnvironmentInvalidRequest("invalid_command")), ); - return yield* orchestrationEngine.dispatch(normalizedCommand).pipe( + const result = yield* orchestrationEngine.dispatch(normalizedCommand).pipe( Effect.tapError(() => cleanupFailedUploadedAttachments(args.payload, normalizedCommand), ), @@ -104,6 +114,11 @@ export const orchestrationHttpApiLayer = HttpApiBuilder.group( failEnvironmentInternal("orchestration_dispatch_failed", cause), ), ); + yield* ProjectCloneTracker.discardCloneForDeletedProject( + projectCloneTracker, + normalizedCommand, + ); + return result; }), ); }), diff --git a/apps/server/src/project/ProjectCloneTracker.test.ts b/apps/server/src/project/ProjectCloneTracker.test.ts new file mode 100644 index 000000000000..d26a2ed712ac --- /dev/null +++ b/apps/server/src/project/ProjectCloneTracker.test.ts @@ -0,0 +1,303 @@ +import { describe, expect, it } from "@effect/vitest"; +import { + OrchestrationDispatchCommandError, + ProjectId, + SourceControlRepositoryError, +} from "@t3tools/contracts"; +import * as Deferred from "effect/Deferred"; +import * as Effect from "effect/Effect"; +import * as Fiber from "effect/Fiber"; +import * as Layer from "effect/Layer"; +import * as Stream from "effect/Stream"; +import * as TestClock from "effect/testing/TestClock"; + +import * as SourceControlRepositoryService from "../sourceControl/SourceControlRepositoryService.ts"; +import * as ProjectCloneTracker from "./ProjectCloneTracker.ts"; +import { parseGitCloneProgressLine } from "./gitCloneProgress.ts"; + +const projectId = ProjectId.make("project-1"); +const startInput = { + projectId, + title: "t3code", + createdAt: "2026-01-01T00:00:00.000Z", + remoteUrl: "git@github.com:octocat/t3code.git", + destinationPath: "/workspace/t3code", +}; + +function makeHarness(options?: { + readonly clone?: SourceControlRepositoryService.SourceControlRepositoryService["Service"]["cloneRepository"]; +}) { + const created: Array<{ projectId: ProjectId; workspaceRoot: string }> = []; + const cloned: Array = []; + const discarded: Array = []; + const hooks: ProjectCloneTracker.ProjectCloneHooks = { + createProject: (input) => + Effect.sync(() => { + created.push({ projectId: input.projectId, workspaceRoot: input.workspaceRoot }); + }), + onCloned: (input) => Effect.sync(() => void cloned.push(input.projectId)), + }; + const layer = ProjectCloneTracker.layer.pipe( + Layer.provide( + Layer.mock(SourceControlRepositoryService.SourceControlRepositoryService)({ + prepareClone: (input) => + Effect.succeed({ + destinationPath: input.destinationPath, + remoteUrl: input.remoteUrl ?? "", + cloneUrl: input.remoteUrl ?? "", + repository: null, + }), + cloneRepository: + options?.clone ?? + ((input) => + Effect.succeed({ + cwd: input.destinationPath, + remoteUrl: input.remoteUrl ?? "", + repository: null, + })), + discardClone: (destination) => Effect.sync(() => void discarded.push(destination)), + }), + ), + ); + return { layer, hooks, created, cloned, discarded }; +} + +describe("ProjectCloneTracker", () => { + it.effect("creates the project first and reports the clone through the stream", () => { + const release = Deferred.makeUnsafe(); + const harness = makeHarness({ + clone: (input, options) => + Effect.gen(function* () { + yield* ( + options?.onProgress?.({ stage: "receiving", percent: 40, detail: "1 MiB" }) ?? + Effect.void + ); + yield* Deferred.await(release); + return { cwd: input.destinationPath, remoteUrl: input.remoteUrl ?? "", repository: null }; + }), + }); + return Effect.gen(function* () { + const tracker = yield* ProjectCloneTracker.ProjectCloneTracker; + const collected = yield* tracker.stream.pipe( + Stream.takeUntil((clones) => clones[0]?.phase === "done"), + Stream.runCollect, + Effect.forkChild, + ); + yield* Effect.yieldNow; + + const result = yield* tracker.start(startInput, harness.hooks); + expect(result.cwd).toBe("/workspace/t3code"); + // The project exists before git runs so the draft can open immediately. + expect(harness.created).toEqual([{ projectId, workspaceRoot: "/workspace/t3code" }]); + + yield* Effect.yieldNow; + const running = yield* tracker.get(projectId); + expect(running).toMatchObject({ phase: "running", stage: "receiving", percent: 40 }); + + yield* Deferred.succeed(release, undefined); + const lists = yield* Fiber.join(collected); + const final = lists.at(-1)?.[0]; + expect(final).toMatchObject({ phase: "done", percent: 100 }); + expect(harness.cloned).toEqual([projectId]); + + // Done clones drop out after the grace window so the toast can settle. + yield* TestClock.adjust("31 seconds"); + expect(yield* tracker.get(projectId)).toBeNull(); + }).pipe(Effect.provide(harness.layer)); + }); + + it.effect("keeps a failed clone with git's own explanation and retries it", () => { + let attempts = 0; + const harness = makeHarness({ + clone: (input) => + Effect.suspend(() => { + attempts += 1; + return attempts === 1 + ? Effect.fail( + new SourceControlRepositoryError({ + operation: "cloneRepository", + provider: "unknown", + detail: "fatal: repository not found", + }), + ) + : Effect.succeed({ + cwd: input.destinationPath, + remoteUrl: input.remoteUrl ?? "", + repository: null, + }); + }), + }); + return Effect.gen(function* () { + const tracker = yield* ProjectCloneTracker.ProjectCloneTracker; + yield* tracker.start(startInput, harness.hooks); + yield* Effect.yieldNow; + const failed = yield* tracker.get(projectId); + expect(failed).toMatchObject({ phase: "failed", error: "fatal: repository not found" }); + + expect(yield* tracker.retry(projectId)).toBe(true); + // The partial checkout is cleared so git sees an empty destination. + expect(harness.discarded).toEqual(["/workspace/t3code"]); + yield* Effect.yieldNow; + expect((yield* tracker.get(projectId))?.phase).toBe("done"); + expect(attempts).toBe(2); + }).pipe(Effect.provide(harness.layer)); + }); + + it.effect("cancel interrupts the clone and removes the partial checkout", () => { + const harness = makeHarness({ clone: () => Effect.never }); + return Effect.gen(function* () { + const tracker = yield* ProjectCloneTracker.ProjectCloneTracker; + yield* tracker.start(startInput, harness.hooks); + yield* Effect.yieldNow; + + expect(yield* tracker.cancel(projectId)).toBe(true); + expect((yield* tracker.get(projectId))?.phase).toBe("cancelled"); + expect(harness.discarded).toEqual(["/workspace/t3code"]); + // Nothing left to cancel; retry is what brings it back. + expect(yield* tracker.cancel(projectId)).toBe(false); + expect(yield* tracker.retry(projectId)).toBe(true); + }).pipe(Effect.provide(harness.layer)); + }); + + it.effect("a cancel that lands after git finished keeps the checkout", () => { + const gate = Deferred.makeUnsafe(); + const harness = makeHarness(); + // The clone itself completes instantly; the post-clone hook is what hangs. + const hooks: ProjectCloneTracker.ProjectCloneHooks = { + ...harness.hooks, + onCloned: () => Deferred.await(gate), + }; + return Effect.gen(function* () { + const tracker = yield* ProjectCloneTracker.ProjectCloneTracker; + yield* tracker.start(startInput, hooks); + yield* Effect.yieldNow; + expect((yield* tracker.get(projectId))?.phase).toBe("done"); + expect(yield* tracker.cancel(projectId)).toBe(false); + expect(harness.discarded).toEqual([]); + yield* Deferred.succeed(gate, undefined); + }).pipe(Effect.provide(harness.layer)); + }); + + it.effect("hands git the credential-bearing URL while snapshots carry the redacted one", () => { + const cloneUrls: Array = []; + const harness = makeHarness({ + clone: (input) => + Effect.sync(() => { + cloneUrls.push(input.remoteUrl ?? ""); + return { cwd: input.destinationPath, remoteUrl: "", repository: null }; + }), + }); + const layer = ProjectCloneTracker.layer.pipe( + Layer.provide( + Layer.mock(SourceControlRepositoryService.SourceControlRepositoryService)({ + prepareClone: (input) => + Effect.succeed({ + destinationPath: input.destinationPath, + remoteUrl: "https://github.com/octocat/t3code.git", + cloneUrl: "https://user:s3cret@github.com/octocat/t3code.git", + repository: null, + }), + cloneRepository: (input) => + Effect.sync(() => { + cloneUrls.push(input.remoteUrl ?? ""); + return { cwd: input.destinationPath, remoteUrl: "", repository: null }; + }), + discardClone: () => Effect.void, + }), + ), + ); + return Effect.gen(function* () { + const tracker = yield* ProjectCloneTracker.ProjectCloneTracker; + const result = yield* tracker.start(startInput, harness.hooks); + yield* Effect.yieldNow; + expect(result.remoteUrl).toBe("https://github.com/octocat/t3code.git"); + expect(cloneUrls).toEqual(["https://user:s3cret@github.com/octocat/t3code.git"]); + }).pipe(Effect.provide(layer)); + }); + + it.effect("discard forgets a project's clone when the project is deleted", () => { + const harness = makeHarness({ clone: () => Effect.never }); + return Effect.gen(function* () { + const tracker = yield* ProjectCloneTracker.ProjectCloneTracker; + yield* tracker.start(startInput, harness.hooks); + yield* Effect.yieldNow; + yield* tracker.discard(projectId); + expect(yield* tracker.get(projectId)).toBeNull(); + expect(harness.discarded).toEqual(["/workspace/t3code"]); + }).pipe(Effect.provide(harness.layer)); + }); + + it.effect("releases the claim when project creation fails", () => { + const harness = makeHarness(); + const hooks: ProjectCloneTracker.ProjectCloneHooks = { + ...harness.hooks, + createProject: () => + Effect.fail(new OrchestrationDispatchCommandError({ message: "workspace root exists" })), + }; + return Effect.gen(function* () { + const tracker = yield* ProjectCloneTracker.ProjectCloneTracker; + const error = yield* Effect.flip(tracker.start(startInput, hooks)); + expect(error.message).toContain("workspace root exists"); + expect(yield* tracker.get(projectId)).toBeNull(); + // The destination is free again for a corrected attempt. + yield* tracker.start(startInput, harness.hooks); + yield* Effect.yieldNow; + expect((yield* tracker.get(projectId))?.phase).toBe("done"); + }).pipe(Effect.provide(harness.layer)); + }); + + it.effect("does not create a project when the clone cannot be prepared", () => { + const harness = makeHarness(); + const layer = ProjectCloneTracker.layer.pipe( + Layer.provide( + Layer.mock(SourceControlRepositoryService.SourceControlRepositoryService)({ + prepareClone: () => + Effect.fail( + new SourceControlRepositoryError({ + operation: "cloneRepository", + provider: "unknown", + detail: "Destination path already exists and is not empty.", + }), + ), + }), + ), + ); + return Effect.gen(function* () { + const tracker = yield* ProjectCloneTracker.ProjectCloneTracker; + const error = yield* Effect.flip(tracker.start(startInput, harness.hooks)); + expect(error.message).toContain("not empty"); + expect(harness.created).toEqual([]); + expect(yield* tracker.get(projectId)).toBeNull(); + }).pipe(Effect.provide(layer)); + }); +}); + +describe("parseGitCloneProgressLine", () => { + it("parses git's transfer counters and ignores other output", () => { + expect( + parseGitCloneProgressLine("Receiving objects: 45% (4500/10000), 12.30 MiB | 5.00 MiB/s"), + ).toEqual({ stage: "receiving", percent: 45, detail: "12.30 MiB | 5.00 MiB/s" }); + expect(parseGitCloneProgressLine("Resolving deltas: 100% (700/700), done.")).toEqual({ + stage: "resolving", + percent: 100, + detail: null, + }); + expect(parseGitCloneProgressLine("remote: Compressing objects: 12% (3/25)")).toEqual({ + stage: "counting", + percent: 12, + detail: null, + }); + expect(parseGitCloneProgressLine("Updating files: 78% (2104/2700)")).toEqual({ + stage: "checkout", + percent: 78, + detail: null, + }); + expect(parseGitCloneProgressLine("remote: Enumerating objects: 10, done.")).toEqual({ + stage: "counting", + percent: null, + detail: null, + }); + expect(parseGitCloneProgressLine("Cloning into 't3code'...")).toBeNull(); + expect(parseGitCloneProgressLine("fatal: repository not found")).toBeNull(); + }); +}); diff --git a/apps/server/src/project/ProjectCloneTracker.ts b/apps/server/src/project/ProjectCloneTracker.ts new file mode 100644 index 000000000000..854c97817529 --- /dev/null +++ b/apps/server/src/project/ProjectCloneTracker.ts @@ -0,0 +1,487 @@ +import type { + OrchestrationCommand, + ProjectCloneSnapshot, + ProjectCloneStage, + ProjectCloneStartInput, + ProjectCloneStartResult, + ProjectId, + SourceControlRepositoryInfo, +} from "@t3tools/contracts"; +import { + OrchestrationDispatchCommandError, + PROJECT_CLONE_DETAIL_MAX_LENGTH, + PROJECT_CLONE_ERROR_MAX_LENGTH, + SourceControlRepositoryError, +} from "@t3tools/contracts"; +import * as Cause from "effect/Cause"; +import * as Context from "effect/Context"; +import * as DateTime from "effect/DateTime"; +import * as Effect from "effect/Effect"; +import * as Exit from "effect/Exit"; +import * as Fiber from "effect/Fiber"; +import * as Layer from "effect/Layer"; +import * as PubSub from "effect/PubSub"; +import * as Queue from "effect/Queue"; +import * as Ref from "effect/Ref"; +import * as Schema from "effect/Schema"; +import * as Semaphore from "effect/Semaphore"; +import * as Scope from "effect/Scope"; +import * as Stream from "effect/Stream"; + +import * as SourceControlRepositoryService from "../sourceControl/SourceControlRepositoryService.ts"; + +/** + * Runs repository clones that back newly added projects and tracks their + * progress so clients can show it anywhere, not just in the surface that + * started the clone. + * + * The clone is detached from the request that started it: the palette closes + * immediately, the project already exists (pointing at the empty destination), + * and the composer can hold a draft for it. Snapshots are memory only. A done + * clone is dropped after a short grace window; a failed one stays until it is + * retried or the server restarts, since the empty project is the durable + * record the user can act on. + */ +export class ProjectCloneTracker extends Context.Service< + ProjectCloneTracker, + { + /** + * Resolves the remote, creates the project, and starts the clone in the + * background. Fails before creating anything when the destination or + * repository is unusable so the caller can report it inline. + */ + readonly start: ( + input: ProjectCloneStartInput, + hooks: ProjectCloneHooks, + ) => Effect.Effect< + ProjectCloneStartResult, + SourceControlRepositoryError | OrchestrationDispatchCommandError + >; + /** Interrupts a running clone and deletes the partial checkout. */ + readonly cancel: (projectId: ProjectId) => Effect.Effect; + /** Restarts a failed or cancelled clone into the same destination. */ + readonly retry: (projectId: ProjectId) => Effect.Effect; + /** + * Stops tracking a project's clone, interrupting it if it still runs and + * removing an unfinished checkout. Called when the project is deleted. + */ + readonly discard: (projectId: ProjectId) => Effect.Effect; + readonly get: (projectId: ProjectId) => Effect.Effect; + /** Emits every tracked clone first, then the full list after each change. */ + readonly stream: Stream.Stream>; + } +>()("t3/project/ProjectCloneTracker") {} + +/** + * Orchestration side effects the caller owns. The tracker never depends on the + * engine directly: the WebSocket handler dispatches with the client's origin, + * and the tracker only needs the outcome. + */ +export interface ProjectCloneHooks { + readonly createProject: (input: { + readonly projectId: ProjectId; + readonly title: string; + readonly workspaceRoot: string; + readonly createdAt: string; + }) => Effect.Effect; + /** Runs after a successful clone so cached repository identity and git status refresh. */ + readonly onCloned: (input: { + readonly projectId: ProjectId; + readonly workspaceRoot: string; + }) => Effect.Effect; +} + +/** Finished snapshots stay visible this long so a late subscriber sees the outcome. */ +const DONE_RETENTION = "30 seconds"; + +function clampText(text: string, maxLength: number): string { + return text.length <= maxLength ? text : `${text.slice(0, maxLength - 1)}…`; +} + +const nowIso = Effect.map(DateTime.now, DateTime.formatIso); + +interface TrackedClone { + readonly snapshot: ProjectCloneSnapshot; + readonly fiber: Fiber.Fiber | null; + readonly hooks: ProjectCloneHooks; + readonly input: { + /** What git is given; may carry credentials and never leaves the server. */ + readonly cloneUrl: string; + readonly destinationPath: string; + readonly repository: SourceControlRepositoryInfo | null; + }; +} + +/** @public Service construction is part of the canonical Effect module API. */ +export const make = Effect.gen(function* () { + const repositories = yield* SourceControlRepositoryService.SourceControlRepositoryService; + const clones = yield* Ref.make(new Map()); + const changes = yield* PubSub.unbounded>(); + const retentionFibers = new Map>(); + let sequence = 0; + // Clone fibers outlive the RPC that started them but not the server. + const cloneScope = yield* Scope.make("parallel"); + yield* Effect.addFinalizer(() => Scope.close(cloneScope, Exit.void)); + // start/cancel/retry/discard mutate the same entry and the same directory; + // one at a time keeps a double-clicked Retry from racing two clones into it. + const actionLock = yield* Semaphore.make(1); + const locked = (effect: Effect.Effect) => actionLock.withPermits(1)(effect); + + const list = Ref.get(clones).pipe( + Effect.map((current) => Array.from(current.values(), (tracked) => tracked.snapshot)), + ); + const publish = list.pipe(Effect.flatMap((snapshots) => PubSub.publish(changes, snapshots))); + + const modify = ( + projectId: ProjectId, + mutate: (tracked: TrackedClone) => TrackedClone, + ): Effect.Effect => + Ref.modify(clones, (current) => { + const existing = current.get(projectId); + if (!existing) return [null, current] as const; + const nextTracked = mutate(existing); + const nextSnapshot = { ...nextTracked.snapshot, sequence: ++sequence }; + const next = new Map(current); + next.set(projectId, { ...nextTracked, snapshot: nextSnapshot }); + return [nextSnapshot, next] as const; + }).pipe(Effect.tap((snapshot) => (snapshot ? publish : Effect.void))); + + const clearRetention = (projectId: ProjectId) => { + const fiber = retentionFibers.get(projectId); + retentionFibers.delete(projectId); + return fiber ? Fiber.interrupt(fiber).pipe(Effect.ignore) : Effect.void; + }; + + const remove = (projectId: ProjectId) => + Ref.update(clones, (current) => { + if (!current.has(projectId)) return current; + const next = new Map(current); + next.delete(projectId); + return next; + }).pipe(Effect.andThen(publish)); + + const scheduleRemoval = (projectId: ProjectId) => + Effect.gen(function* () { + yield* clearRetention(projectId); + const fiber = yield* remove(projectId).pipe( + Effect.delay(DONE_RETENTION), + Effect.ensuring( + Effect.sync(() => { + if (retentionFibers.get(projectId) === fiber) retentionFibers.delete(projectId); + }), + ), + Effect.forkDetach, + ); + retentionFibers.set(projectId, fiber); + }); + + const progress = ( + projectId: ProjectId, + update: { + readonly stage: ProjectCloneStage; + readonly percent: number | null; + readonly detail: string | null; + }, + ) => + modify(projectId, (tracked) => ({ + ...tracked, + snapshot: { + ...tracked.snapshot, + stage: update.stage, + percent: update.percent, + detail: + update.detail === null ? null : clampText(update.detail, PROJECT_CLONE_DETAIL_MAX_LENGTH), + }, + })).pipe(Effect.asVoid); + + const finish = ( + projectId: ProjectId, + phase: "done" | "failed" | "cancelled", + error: string | null, + ) => + Effect.gen(function* () { + const endedAt = yield* nowIso; + yield* modify(projectId, (tracked) => ({ + ...tracked, + fiber: null, + snapshot: { + ...tracked.snapshot, + phase, + endedAt, + percent: phase === "done" ? 100 : tracked.snapshot.percent, + error: error === null ? null : clampText(error, PROJECT_CLONE_ERROR_MAX_LENGTH), + }, + })); + if (phase === "done") yield* scheduleRemoval(projectId); + }); + + /** + * The clone body. Runs in its own fiber; the tracker records the outcome + * through `onExit`, which still runs when the fiber is interrupted (a + * `matchCause` handler would be skipped). Once git has finished the clone + * is marked done before the post-clone hook runs, so a late Cancel cannot + * tear down a complete checkout. + */ + const runClone = (projectId: ProjectId, tracked: TrackedClone) => + repositories + .cloneRepository( + { remoteUrl: tracked.input.cloneUrl, destinationPath: tracked.input.destinationPath }, + { onProgress: (update) => progress(projectId, update), timeoutMs: null }, + ) + .pipe( + Effect.onExit((exit) => + Exit.isSuccess(exit) + ? finish(projectId, "done", null) + : Cause.hasInterruptsOnly(exit.cause) + ? finish(projectId, "cancelled", null) + : finish(projectId, "failed", describeCloneFailure(exit.cause)), + ), + Effect.flatMap(() => + tracked.hooks + .onCloned({ projectId, workspaceRoot: tracked.input.destinationPath }) + .pipe(Effect.ignoreCause({ log: true })), + ), + Effect.ignoreCause(), + ); + + const launch = (projectId: ProjectId) => + Effect.gen(function* () { + const current = yield* Ref.get(clones); + const tracked = current.get(projectId); + if (!tracked) return; + const fiber = yield* runClone(projectId, tracked).pipe(Effect.forkIn(cloneScope)); + yield* Ref.update(clones, (map) => { + const existing = map.get(projectId); + if (!existing) return map; + const next = new Map(map); + next.set(projectId, { ...existing, fiber }); + return next; + }); + }); + + const start: ProjectCloneTracker["Service"]["start"] = Effect.fn("ProjectCloneTracker.start")( + function* (input, hooks) { + const prepared = yield* repositories.prepareClone(input); + const startedAt = yield* nowIso; + const snapshot: ProjectCloneSnapshot = { + projectId: input.projectId, + remoteUrl: prepared.remoteUrl, + destinationPath: prepared.destinationPath, + repository: prepared.repository, + phase: "running", + stage: "connecting", + percent: null, + detail: null, + error: null, + startedAt, + endedAt: null, + sequence: ++sequence, + }; + const claimed = yield* Ref.modify(clones, (current) => { + // A second start for the same project, or for a destination another + // clone already owns, must not race two gits into one directory. + const conflict = Array.from(current.values()).some( + (tracked) => + tracked.snapshot.projectId === input.projectId || + tracked.input.destinationPath === prepared.destinationPath, + ); + if (conflict) return [false, current] as const; + const next = new Map(current); + next.set(input.projectId, { + snapshot, + fiber: null, + hooks, + input: { + cloneUrl: prepared.cloneUrl, + destinationPath: prepared.destinationPath, + repository: prepared.repository, + }, + }); + return [true, next] as const; + }); + if (!claimed) { + return yield* new SourceControlRepositoryError({ + operation: "cloneRepository", + provider: input.provider ?? "unknown", + detail: "A clone into this destination is already in progress.", + }); + } + // Everything after the claim runs to completion even if the requesting + // connection drops: a claimed entry with no fiber could neither be + // cancelled nor retried. Any failure in here releases the claim. + yield* Effect.uninterruptible( + Effect.gen(function* () { + yield* clearRetention(input.projectId); + // The entry is registered before the project exists so a + // thread.create or project.delete racing this call already sees it. + yield* hooks.createProject({ + projectId: input.projectId, + title: input.title, + workspaceRoot: prepared.destinationPath, + createdAt: input.createdAt, + }); + yield* publish; + yield* launch(input.projectId); + }).pipe(Effect.tapError(() => remove(input.projectId))), + ); + return { + projectId: input.projectId, + cwd: prepared.destinationPath, + remoteUrl: prepared.remoteUrl, + repository: prepared.repository, + }; + }, + ); + + const get: ProjectCloneTracker["Service"]["get"] = (projectId) => + Ref.get(clones).pipe(Effect.map((current) => current.get(projectId)?.snapshot ?? null)); + + const cancel: ProjectCloneTracker["Service"]["cancel"] = (projectId) => + Effect.gen(function* () { + const current = yield* Ref.get(clones); + const tracked = current.get(projectId); + if (!tracked || tracked.snapshot.phase !== "running" || !tracked.fiber) return false; + // Uninterruptible past this point: a client that disconnects mid-cancel + // must not leave a dead fiber behind a "running" snapshot. + yield* Effect.uninterruptible( + Effect.gen(function* () { + yield* Fiber.interrupt(tracked.fiber!); + const after = yield* get(projectId); + // Git finished in the window before the interrupt landed: keep it. + if (after?.phase === "done") return; + if (after?.phase === "running") yield* finish(projectId, "cancelled", null); + // The partial checkout goes so a retry starts from an empty destination. + yield* repositories.discardClone(tracked.input.destinationPath).pipe(Effect.ignore); + }), + ); + return true; + }); + + const retry: ProjectCloneTracker["Service"]["retry"] = (projectId) => + Effect.gen(function* () { + const current = yield* Ref.get(clones); + const tracked = current.get(projectId); + if ( + !tracked || + (tracked.snapshot.phase !== "failed" && tracked.snapshot.phase !== "cancelled") + ) { + return false; + } + // A retry into leftover files would fail on the non-empty destination + // with a less useful message, so a cleanup failure is the error here. + yield* repositories.discardClone(tracked.input.destinationPath); + const startedAt = yield* nowIso; + yield* modify(projectId, (entry) => ({ + ...entry, + snapshot: { + ...entry.snapshot, + phase: "running", + stage: "connecting", + percent: null, + detail: null, + error: null, + startedAt, + endedAt: null, + }, + })); + yield* launch(projectId); + return true; + }); + + const discard: ProjectCloneTracker["Service"]["discard"] = (projectId) => + Effect.gen(function* () { + const current = yield* Ref.get(clones); + const tracked = current.get(projectId); + if (!tracked) return; + if (tracked.fiber) yield* Fiber.interrupt(tracked.fiber); + // Git may have finished while the interrupt was landing; the project + // is going away either way, but a complete checkout is the user's. + const after = yield* get(projectId); + if (after?.phase !== "done") { + yield* repositories.discardClone(tracked.input.destinationPath).pipe(Effect.ignore); + } + yield* clearRetention(projectId); + yield* remove(projectId); + }); + + // One-slot sliding mailbox per subscriber: a slow socket only ever holds the + // newest list, and lists are whole states so skipping intermediates is safe. + const stream: ProjectCloneTracker["Service"]["stream"] = Stream.callback< + ReadonlyArray + >( + (mailbox) => + Effect.gen(function* () { + const subscription = yield* PubSub.subscribe(changes); + Queue.offerUnsafe(mailbox, yield* list); + yield* Stream.fromSubscription(subscription).pipe( + Stream.runForEach((snapshots) => + Effect.sync(() => Queue.offerUnsafe(mailbox, snapshots)), + ), + Effect.forkScoped, + ); + }), + { bufferSize: 1, strategy: "sliding" }, + ); + + return ProjectCloneTracker.of({ + start: (input, hooks) => locked(start(input, hooks)), + cancel: (projectId) => locked(cancel(projectId)), + retry: (projectId) => locked(retry(projectId)), + discard: (projectId) => locked(discard(projectId)), + get, + stream, + }); +}); + +const isSourceControlRepositoryError = Schema.is(SourceControlRepositoryError); + +/** + * A project whose clone has not landed has no files to work in. Every + * dispatch transport (WebSocket, HTTP) runs this before normalizing so no + * client can start a thread on an empty tree, and no attachment copies are + * made for a command that is about to be refused. + */ +export const rejectCommandsDuringClone = ( + tracker: ProjectCloneTracker["Service"], + command: { readonly type: string; readonly projectId?: ProjectId; readonly bootstrap?: unknown }, +): Effect.Effect => + Effect.gen(function* () { + const projectId = + command.type === "thread.create" + ? (command.projectId ?? null) + : command.type === "thread.turn.start" + ? bootstrapProjectId(command.bootstrap) + : null; + if (projectId === null) return; + const clone = yield* tracker.get(projectId); + if (clone === null || clone.phase === "done") return; + return yield* new OrchestrationDispatchCommandError({ + message: + clone.phase === "running" + ? "The repository is still being cloned." + : "The repository was not cloned. Retry the clone first.", + }); + }); + +function bootstrapProjectId(bootstrap: unknown): ProjectId | null { + if (typeof bootstrap !== "object" || bootstrap === null) return null; + const createThread = (bootstrap as { createThread?: { projectId?: ProjectId } }).createThread; + return createThread?.projectId ?? null; +} + +/** Removing a project mid-clone stops the clone and drops its partial checkout. */ +export const discardCloneForDeletedProject = ( + tracker: ProjectCloneTracker["Service"], + command: OrchestrationCommand, +): Effect.Effect => + command.type === "project.delete" ? tracker.discard(command.projectId) : Effect.void; + +function describeCloneFailure(cause: Cause.Cause): string { + const error = Cause.squash(cause); + if (isSourceControlRepositoryError(error)) return error.detail; + return error instanceof Error && error.message.trim().length > 0 + ? error.message + : "The repository could not be cloned."; +} + +export const layer = Layer.effect(ProjectCloneTracker, make); diff --git a/apps/server/src/project/gitCloneProgress.ts b/apps/server/src/project/gitCloneProgress.ts new file mode 100644 index 000000000000..a6108442b85a --- /dev/null +++ b/apps/server/src/project/gitCloneProgress.ts @@ -0,0 +1,44 @@ +import type { ProjectCloneStage } from "@t3tools/contracts"; + +export interface GitCloneProgressLine { + readonly stage: ProjectCloneStage; + readonly percent: number | null; + /** Transfer detail after the count, e.g. `12.30 MiB | 5.00 MiB/s`. */ + readonly detail: string | null; +} + +const STAGE_PREFIXES: ReadonlyArray = [ + [/^remote: Enumerating objects/, "counting"], + [/^remote: Counting objects/, "counting"], + [/^remote: Compressing objects/, "counting"], + [/^Receiving objects/, "receiving"], + [/^Resolving deltas/, "resolving"], + [/^Updating files/, "checkout"], + [/^Checking out files/, "checkout"], +]; + +const PERCENT = /:\s+(\d+)%\s+\((\d+)\/(\d+)\)(?:,\s*(.*?))?\s*(?:,\s*done\.)?\s*$/; + +/** + * Parses one line of `git clone --progress` stderr. Git redraws each counter + * with a bare `\r`, so callers hand over each redraw as its own line. Lines + * that are not progress counters (hints, warnings, `Cloning into ...`) return + * null and are left for the error surface. + */ +export function parseGitCloneProgressLine(line: string): GitCloneProgressLine | null { + const trimmed = line.trim(); + const stageEntry = STAGE_PREFIXES.find(([pattern]) => pattern.test(trimmed)); + if (!stageEntry) return null; + const stage = stageEntry[1]; + const match = PERCENT.exec(trimmed); + if (!match) return { stage, percent: null, detail: null }; + const percent = Number(match[1]); + const rawDetail = match[4]?.trim() ?? ""; + // The trailer of a finished line is "done." which carries no information. + const detail = rawDetail.length > 0 && rawDetail !== "done." ? rawDetail : null; + return { + stage, + percent: Number.isFinite(percent) ? Math.max(0, Math.min(100, percent)) : null, + detail, + }; +} diff --git a/apps/server/src/server.test.ts b/apps/server/src/server.test.ts index d5ba7bc6bf56..7402e908fcf8 100644 --- a/apps/server/src/server.test.ts +++ b/apps/server/src/server.test.ts @@ -148,6 +148,7 @@ import * as ServerRuntimeStartup from "./serverRuntimeStartup.ts"; import * as ServiceLauncherClient from "./cloud/serviceLauncherClient.ts"; import * as ServerSettings from "./serverSettings.ts"; import * as TerminalManager from "./terminal/Manager.ts"; +import * as ProjectCloneTracker from "./project/ProjectCloneTracker.ts"; import * as WorktreeSetupTracker from "./project/WorktreeSetupTracker.ts"; import * as PreviewManager from "./preview/Manager.ts"; import * as PortScanner from "./preview/PortScanner.ts"; @@ -933,6 +934,13 @@ const buildAppUnderTest = (options?: { ...options?.layers?.terminalManager, }), WorktreeSetupTracker.layer, + ProjectCloneTracker.layer.pipe( + Layer.provide( + Layer.mock(SourceControlRepositoryService.SourceControlRepositoryService)({ + ...options?.layers?.sourceControlRepositoryService, + }), + ), + ), ), ), Layer.provide( @@ -7238,6 +7246,106 @@ it.layer(NodeServices.layer)("server router seam", (it) => { }).pipe(Effect.provide(NodeHttpServer.layerTest)), ); + it.effect("starts a project clone in the background and blocks threads until it lands", () => + Effect.gen(function* () { + const fs = yield* FileSystem.FileSystem; + const path = yield* Path.Path; + const parentDir = yield* fs.makeTempDirectoryScoped({ prefix: "t3-ws-project-clone-" }); + const destinationPath = path.join(parentDir, "t3code"); + const projectId = ProjectId.make("project-clone-1"); + const dispatched: Array = []; + const cloneGate = yield* Deferred.make(); + const metaUpdateDispatched = yield* Deferred.make(); + + yield* buildAppUnderTest({ + layers: { + orchestrationEngine: { + dispatch: (command) => + Effect.sync(() => { + dispatched.push(command.type); + return { sequence: dispatched.length }; + }).pipe( + Effect.tap(() => + command.type === "project.meta.update" + ? Deferred.succeed(metaUpdateDispatched, undefined) + : Effect.void, + ), + ), + }, + sourceControlRepositoryService: { + prepareClone: (input) => + Effect.succeed({ + destinationPath: input.destinationPath, + remoteUrl: input.remoteUrl ?? "", + cloneUrl: input.remoteUrl ?? "", + repository: null, + }), + cloneRepository: (input) => + Deferred.await(cloneGate).pipe( + Effect.as({ + cwd: input.destinationPath, + remoteUrl: input.remoteUrl ?? "", + repository: null, + }), + ), + }, + }, + }); + + const wsUrl = yield* getWsServerUrl("/ws"); + yield* Effect.scoped( + withWsRpcClient(wsUrl, (client) => + Effect.gen(function* () { + const started = yield* client[WS_METHODS.projectCloneStart]({ + projectId, + title: "t3code", + createdAt: "2026-01-01T00:00:00.000Z", + remoteUrl: "git@github.com:octocat/t3code.git", + destinationPath, + }); + assert.equal(started.cwd, destinationPath); + // The project exists before the clone finishes. + assert.deepEqual(dispatched, ["project.create"]); + + const blocked = yield* Effect.flip( + client[ORCHESTRATION_WS_METHODS.dispatchCommand]({ + type: "thread.create", + commandId: CommandId.make("cmd-thread-create-while-cloning"), + threadId: ThreadId.make("thread-while-cloning"), + projectId, + title: "Draft", + modelSelection: { + instanceId: ProviderInstanceId.make("codex"), + model: "gpt-5-codex", + }, + runtimeMode: "full-access", + interactionMode: "default", + branch: null, + worktreePath: null, + createdAt: "2026-01-01T00:00:01.000Z", + }), + ); + assert.include(String(blocked.message), "still being cloned"); + + const snapshots = yield* client[WS_METHODS.subscribeProjectClones]({}).pipe( + Stream.takeUntil((clones) => clones[0]?.phase === "done"), + Stream.runCollect, + Effect.forkChild, + ); + yield* Effect.yieldNow; + yield* Deferred.succeed(cloneGate, undefined); + const lists = yield* Fiber.join(snapshots); + assert.equal(lists.at(-1)?.[0]?.phase, "done"); + // The finished clone refreshes the project so its repository + // identity updates. That hook runs after the done snapshot. + yield* Deferred.await(metaUpdateDispatched); + assert.deepEqual(dispatched, ["project.create", "project.meta.update"]); + }), + ), + ); + }).pipe(Effect.provide(NodeHttpServer.layerTest)), + ); + it.effect("records thread analytics only after a client command succeeds", () => Effect.gen(function* () { const effects: string[] = []; diff --git a/apps/server/src/server.ts b/apps/server/src/server.ts index 189d62dd8362..cd867c0e4651 100644 --- a/apps/server/src/server.ts +++ b/apps/server/src/server.ts @@ -103,6 +103,7 @@ import * as VcsProjectConfig from "./vcs/VcsProjectConfig.ts"; import * as VcsProcess from "./vcs/VcsProcess.ts"; import * as VcsProvisioningService from "./vcs/VcsProvisioningService.ts"; import * as VcsStatusBroadcaster from "./vcs/VcsStatusBroadcaster.ts"; +import * as ProjectCloneTracker from "./project/ProjectCloneTracker.ts"; import * as GitWorkflowService from "./git/GitWorkflowService.ts"; import * as ReviewService from "./review/ReviewService.ts"; import * as SourceControlProviderRegistry from "./sourceControl/SourceControlProviderRegistry.ts"; @@ -354,6 +355,10 @@ const SourceControlRepositoryServiceLayerLive = SourceControlRepositoryService.l Layer.provideMerge(SourceControlProviderRegistryLayerLive), ); +const ProjectCloneTrackerLayerLive = ProjectCloneTracker.layer.pipe( + Layer.provide(SourceControlRepositoryServiceLayerLive), +); + const ReviewLayerLive = ReviewService.layer.pipe( Layer.provideMerge(GitVcsDriver.layer), Layer.provideMerge(VcsDriverRegistryLayerLive), @@ -366,6 +371,7 @@ const VcsLayerLive = Layer.empty.pipe( Layer.provideMerge(GitWorkflowLayerLive), Layer.provideMerge(ReviewLayerLive), Layer.provideMerge(SourceControlRepositoryServiceLayerLive), + Layer.provideMerge(ProjectCloneTrackerLayerLive), Layer.provideMerge( VcsStatusBroadcaster.layer.pipe( Layer.provide(GitWorkflowLayerLive), diff --git a/apps/server/src/sourceControl/SourceControlRepositoryService.test.ts b/apps/server/src/sourceControl/SourceControlRepositoryService.test.ts index 461bff08668a..da45f9eabf2b 100644 --- a/apps/server/src/sourceControl/SourceControlRepositoryService.test.ts +++ b/apps/server/src/sourceControl/SourceControlRepositoryService.test.ts @@ -178,7 +178,7 @@ it.effect("clones a looked-up repository into the requested destination", () => assert.deepStrictEqual(cloneCalls, [ { cwd: parent, - args: ["clone", CLONE_URLS.url, "t3code"], + args: ["clone", "--progress", CLONE_URLS.url, "t3code"], }, ]); }).pipe( @@ -197,6 +197,154 @@ it.effect("clones a looked-up repository into the requested destination", () => }).pipe(Effect.provide(NodeServices.layer)), ); +it.effect("reports clone progress from git's stderr and keeps its error text on failure", () => + Effect.gen(function* () { + const fs = yield* FileSystem.FileSystem; + const path = yield* Path.Path; + const parent = yield* fs.makeTempDirectoryScoped({ + prefix: "t3-source-control-clone-progress-", + }); + const destinationPath = path.join(parent, "t3code"); + const progress: Array<{ stage: string; percent: number | null; detail: string | null }> = []; + + const stderrLines = [ + "Cloning into 't3code'...", + "remote: Enumerating objects: 10, done.", + "Receiving objects: 40% (4/10), 1.00 MiB | 2.00 MiB/s", + "Receiving objects: 100% (10/10), 2.50 MiB | 2.00 MiB/s, done.", + "fatal: early EOF", + "fatal: unable to access 'https://user:s3c@ret@github.com/octocat/t3code.git/': could not resolve host", + ]; + const error = yield* Effect.gen(function* () { + const service = yield* SourceControlRepositoryService.SourceControlRepositoryService; + return yield* Effect.flip( + service.cloneRepository( + { remoteUrl: CLONE_URLS.sshUrl, destinationPath }, + { onProgress: (line) => Effect.sync(() => void progress.push(line)) }, + ), + ); + }).pipe( + Effect.provide( + makeLayer({ + git: { + execute: (input) => + Effect.gen(function* () { + for (const line of stderrLines) { + yield* input.progress?.onStderrLine?.(line) ?? Effect.void; + } + return yield* new GitCommandError({ + operation: input.operation, + command: "git", + cwd: input.cwd, + detail: "Git command exited with a non-zero status.", + exitCode: 128, + }); + }), + }, + }), + ), + ); + + assert.deepStrictEqual(progress, [ + { stage: "counting", percent: null, detail: null }, + { stage: "receiving", percent: 40, detail: "1.00 MiB | 2.00 MiB/s" }, + { stage: "receiving", percent: 100, detail: "2.50 MiB | 2.00 MiB/s" }, + ]); + // Git echoes the remote in some failures; the credentials must not follow. + assert.strictEqual( + error.detail, + "fatal: early EOF fatal: unable to access 'https://github.com/octocat/t3code.git/': could not resolve host", + ); + }).pipe(Effect.provide(NodeServices.layer)), +); + +it.effect("strips embedded credentials from the remote URL it reports", () => + Effect.gen(function* () { + const fs = yield* FileSystem.FileSystem; + const path = yield* Path.Path; + const parent = yield* fs.makeTempDirectoryScoped({ prefix: "t3-source-control-redact-" }); + const destinationPath = path.join(parent, "t3code"); + const cloneArgs: Array> = []; + const result = yield* Effect.gen(function* () { + const service = yield* SourceControlRepositoryService.SourceControlRepositoryService; + return yield* service.prepareClone({ + remoteUrl: "https://user:s3cret@github.com/octocat/t3code.git", + destinationPath, + }); + }).pipe( + Effect.provide( + makeLayer({ + git: { + execute: (input) => + Effect.sync(() => { + cloneArgs.push(input.args); + return processOutput(); + }), + }, + }), + ), + ); + assert.equal(result.remoteUrl, "https://github.com/octocat/t3code.git"); + // Git itself still receives the credentials. + assert.equal(result.cloneUrl, "https://user:s3cret@github.com/octocat/t3code.git"); + }).pipe(Effect.provide(NodeServices.layer)), +); + +it.effect("discards only a directory git wrote to", () => + Effect.gen(function* () { + const fs = yield* FileSystem.FileSystem; + const path = yield* Path.Path; + const parent = yield* fs.makeTempDirectoryScoped({ prefix: "t3-source-control-discard-" }); + const partial = path.join(parent, "partial"); + yield* fs.makeDirectory(path.join(partial, ".git"), { recursive: true }); + yield* fs.writeFileString(path.join(partial, "README.md"), "half"); + const foreign = path.join(parent, "foreign"); + yield* fs.makeDirectory(foreign); + yield* fs.writeFileString(path.join(foreign, "notes.txt"), "mine"); + + // A file where the directory should be must not be removed either. + const replaced = path.join(parent, "replaced"); + yield* fs.writeFileString(replaced, "not a directory"); + + yield* Effect.gen(function* () { + const service = yield* SourceControlRepositoryService.SourceControlRepositoryService; + yield* service.discardClone(partial); + const error = yield* Effect.flip(service.discardClone(foreign)); + assert.include(error.detail, "not from the clone"); + const replacedError = yield* Effect.flip(service.discardClone(replaced)); + assert.include(replacedError.detail, "could not be inspected"); + // A destination that never got created is nothing to discard. + yield* service.discardClone(path.join(parent, "missing")); + }).pipe(Effect.provide(makeLayer({}))); + + // The partial clone is emptied but its directory (the workspace root) stays. + assert.deepStrictEqual(yield* fs.readDirectory(partial), []); + assert.deepStrictEqual(yield* fs.readDirectory(foreign), ["notes.txt"]); + assert.strictEqual(yield* fs.readFileString(replaced), "not a directory"); + }).pipe(Effect.provide(NodeServices.layer)), +); + +it.effect("redacts query tokens and userinfo containing '@' from reported URLs", () => + Effect.gen(function* () { + const fs = yield* FileSystem.FileSystem; + const path = yield* Path.Path; + const parent = yield* fs.makeTempDirectoryScoped({ prefix: "t3-source-control-redact2-" }); + yield* Effect.gen(function* () { + const service = yield* SourceControlRepositoryService.SourceControlRepositoryService; + const query = yield* service.prepareClone({ + remoteUrl: "https://github.com/octocat/t3code.git?access_token=s3cret", + destinationPath: path.join(parent, "a"), + }); + assert.equal(query.remoteUrl, "https://github.com/octocat/t3code.git"); + const nested = yield* service.prepareClone({ + remoteUrl: "https://user:pa@rt@github.com/octocat/t3code.git", + destinationPath: path.join(parent, "b"), + }); + assert.equal(nested.remoteUrl, "https://github.com/octocat/t3code.git"); + }).pipe(Effect.provide(makeLayer({}))); + }).pipe(Effect.provide(NodeServices.layer)), +); + it.effect("preserves destination probe failures instead of treating them as missing paths", () => { const fileSystemCause = PlatformError.systemError({ _tag: "PermissionDenied", diff --git a/apps/server/src/sourceControl/SourceControlRepositoryService.ts b/apps/server/src/sourceControl/SourceControlRepositoryService.ts index 0addeca9e785..9bb3b1e029b6 100644 --- a/apps/server/src/sourceControl/SourceControlRepositoryService.ts +++ b/apps/server/src/sourceControl/SourceControlRepositoryService.ts @@ -3,6 +3,7 @@ import * as Effect from "effect/Effect"; import * as FileSystem from "effect/FileSystem"; import * as Layer from "effect/Layer"; import * as Path from "effect/Path"; +import * as Schedule from "effect/Schedule"; import * as Schema from "effect/Schema"; import { @@ -20,6 +21,10 @@ import { import { ServerConfig } from "../config.ts"; import { expandHomePathWith } from "../pathExpansion.ts"; +import { + parseGitCloneProgressLine, + type GitCloneProgressLine, +} from "../project/gitCloneProgress.ts"; import * as GitVcsDriver from "../vcs/GitVcsDriver.ts"; import * as SourceControlProviderRegistry from "./SourceControlProviderRegistry.ts"; const isSourceControlRepositoryError = Schema.is(SourceControlRepositoryError); @@ -30,15 +35,56 @@ export class SourceControlRepositoryService extends Context.Service< readonly lookupRepository: ( input: SourceControlRepositoryLookupInput, ) => Effect.Effect; + /** + * Everything `cloneRepository` checks before running git: the resolved + * remote, the normalized destination, and that the destination is empty. + * Lets a caller create the project first and clone afterwards. + */ + readonly prepareClone: ( + input: SourceControlCloneRepositoryInput, + ) => Effect.Effect; readonly cloneRepository: ( input: SourceControlCloneRepositoryInput, + options?: SourceControlCloneOptions, ) => Effect.Effect; + /** Removes a partial or failed clone so the destination is empty again. */ + readonly discardClone: ( + destinationPath: string, + ) => Effect.Effect; readonly publishRepository: ( input: SourceControlPublishRepositoryInput, ) => Effect.Effect; } >()("t3/sourceControl/SourceControlRepositoryService") {} +export interface SourceControlPreparedClone { + readonly destinationPath: string; + /** Credential-free; safe to show and to store in snapshots. */ + readonly remoteUrl: string; + /** What git is given; may carry embedded credentials. */ + readonly cloneUrl: string; + readonly repository: SourceControlRepositoryInfo | null; +} + +export interface SourceControlCloneOptions { + readonly onProgress?: (line: GitCloneProgressLine) => Effect.Effect; + /** Overrides the default clone budget; `null` disables the deadline. */ + readonly timeoutMs?: number | null; +} + +// The synchronous RPC (older clients, mobile) keeps a deadline: nothing else +// tells the user a clone stalled. The tracked path passes null and relies on +// progress and Cancel instead. +const CLONE_TIMEOUT_MS = 120_000; +const CLONE_ENV = { + // `--progress` forces the transfer counters through the pipe; the delay env + // makes the checkout counter start immediately. No tty means a credential + // prompt would hang forever, so tell git to fail instead. + GIT_PROGRESS_DELAY: "0", + GIT_TERMINAL_PROMPT: "0", + LC_ALL: "C", +} satisfies NodeJS.ProcessEnv; + function mapRepositoryError(operation: string, provider: SourceControlProviderKind) { return Effect.mapError((cause: unknown) => isSourceControlRepositoryError(cause) @@ -64,6 +110,36 @@ function toRepositoryInfo( }; } +/** + * The URL clients see. A pasted `https://user:token@host/…` must not travel + * back over `subscribeProjectClones` to every reader; git still gets the + * original. + */ +function redactRemoteUrl(remoteUrl: string): string { + try { + const url = new URL(remoteUrl); + // Clone URLs have no legitimate query; when one is present it is a token. + if (url.username.length === 0 && url.password.length === 0 && url.search.length === 0) { + return remoteUrl; + } + url.username = ""; + url.password = ""; + url.search = ""; + return url.toString(); + } catch { + return remoteUrl; + } +} + +// Userinfo may itself contain `@`; everything up to the last one before the +// host boundary goes. +const URL_WITH_USERINFO = /\b([a-z][a-z0-9+.-]*:\/\/)[^\s/]+@/gi; + +/** Drops `user:token@` from any URL embedded in free text. */ +function redactUrlCredentials(text: string): string { + return text.replace(URL_WITH_USERINFO, "$1"); +} + function selectRemoteUrl( urls: SourceControlRepositoryCloneUrls, protocol: SourceControlCloneProtocol | undefined, @@ -168,7 +244,7 @@ export const make = Effect.gen(function* () { }, ); - const cloneRepository = Effect.fn("SourceControlRepositoryService.cloneRepository")(function* ( + const prepareClone = Effect.fn("SourceControlRepositoryService.prepareClone")(function* ( input: SourceControlCloneRepositoryInput, ) { const preparedDestination = yield* prepareDestination(input.destinationPath); @@ -194,21 +270,121 @@ export const make = Effect.gen(function* () { }); } - yield* git.execute({ - operation: "SourceControlRepositoryService.cloneRepository", - cwd: preparedDestination.parentPath, - args: ["clone", remoteUrl, preparedDestination.directoryName], - timeoutMs: 120_000, - maxOutputBytes: 256 * 1024, - }); - return { - cwd: preparedDestination.destinationPath, - remoteUrl, + destinationPath: preparedDestination.destinationPath, + remoteUrl: redactRemoteUrl(remoteUrl), + cloneUrl: remoteUrl, repository, + } satisfies SourceControlPreparedClone; + }); + + const cloneRepository = Effect.fn("SourceControlRepositoryService.cloneRepository")(function* ( + input: SourceControlCloneRepositoryInput, + options?: SourceControlCloneOptions, + ) { + const prepared = yield* prepareClone(input); + const onProgress = options?.onProgress; + // Git interleaves progress redraws with its real messages on stderr. The + // last non-progress lines are what explain a failure ("Repository not + // found", "Permission denied"), so keep them for the error detail. + const stderrTail: Array = []; + const onStderrLine = (line: string) => { + const parsed = parseGitCloneProgressLine(line); + if (parsed) return onProgress ? onProgress(parsed) : Effect.void; + return Effect.sync(() => { + const trimmed = line.trim(); + if (trimmed.length === 0 || trimmed.startsWith("Cloning into")) return; + // Git echoes the remote in some failures; the tail becomes user-facing text. + stderrTail.push(redactUrlCredentials(trimmed)); + if (stderrTail.length > 4) stderrTail.shift(); + }); + }; + yield* git + .execute({ + operation: "SourceControlRepositoryService.cloneRepository", + cwd: path.dirname(prepared.destinationPath), + args: ["clone", "--progress", prepared.cloneUrl, path.basename(prepared.destinationPath)], + timeoutMs: options?.timeoutMs === undefined ? CLONE_TIMEOUT_MS : options.timeoutMs, + // Progress redraws add up on a slow multi-GB clone. The buffered copy + // is never read (the tail is kept by hand above), so keep it small + // and let the line callbacks keep flowing past the cap. + maxOutputBytes: 256 * 1024, + appendTruncationMarker: true, + keepLineCallbacksAfterTruncation: true, + env: CLONE_ENV, + progress: { onStderrLine }, + }) + .pipe( + Effect.mapError( + (cause) => + new SourceControlRepositoryError({ + operation: "cloneRepository", + provider: input.provider ?? "unknown", + detail: + stderrTail.length > 0 + ? stderrTail.join(" ") + : "The repository could not be cloned.", + cause, + }), + ), + ); + + return { + cwd: prepared.destinationPath, + remoteUrl: prepared.remoteUrl, + repository: prepared.repository, }; }); + const discardClone = Effect.fn("SourceControlRepositoryService.discardClone")(function* ( + destinationPath: string, + ) { + const normalized = yield* normalizeDestinationPath(destinationPath); + // Only what git left behind may go. The destination was empty when the + // clone started, so anything without a `.git` inside was put there by + // someone else since; refuse rather than delete their files. + // A missing destination is already discarded; any other read failure + // (a file in its place, permissions) is not something to remove through. + const entries = yield* fileSystem.readDirectory(normalized).pipe( + Effect.catchIf( + (cause) => cause.reason._tag === "NotFound", + () => Effect.succeed>([]), + ), + Effect.mapError( + (cause) => + new SourceControlRepositoryError({ + operation: "discardClone", + provider: "unknown", + detail: "The clone destination could not be inspected.", + cause, + }), + ), + ); + if (entries.length > 0 && !entries.includes(".git")) { + return yield* new SourceControlRepositoryError({ + operation: "discardClone", + provider: "unknown", + detail: "Destination path contains files that are not from the clone.", + }); + } + // The directory itself is the project's workspace root and must stay; + // only git's partial contents go. An interrupted git may still be closing + // files, so removal retries briefly. + yield* fileSystem.remove(normalized, { recursive: true, force: true }).pipe( + Effect.andThen(fileSystem.makeDirectory(normalized, { recursive: true })), + Effect.retry({ schedule: Schedule.spaced("200 millis"), times: 5 }), + Effect.mapError( + (cause) => + new SourceControlRepositoryError({ + operation: "discardClone", + provider: "unknown", + detail: "The partial clone could not be removed.", + cause, + }), + ), + ); + }); + const publishRepository = Effect.fn("SourceControlRepositoryService.publishRepository")( function* (input: SourceControlPublishRepositoryInput) { const providerKind = yield* ensureConcreteProvider({ @@ -269,10 +445,14 @@ export const make = Effect.gen(function* () { return SourceControlRepositoryService.of({ lookupRepository: (input) => lookupRepository(input).pipe(mapRepositoryError("lookupRepository", input.provider)), - cloneRepository: (input) => - cloneRepository(input).pipe( + prepareClone: (input) => + prepareClone(input).pipe(mapRepositoryError("cloneRepository", input.provider ?? "unknown")), + cloneRepository: (input, options) => + cloneRepository(input, options).pipe( mapRepositoryError("cloneRepository", input.provider ?? "unknown"), ), + discardClone: (destinationPath) => + discardClone(destinationPath).pipe(mapRepositoryError("discardClone", "unknown")), publishRepository: (input) => publishRepository(input).pipe(mapRepositoryError("publishRepository", input.provider)), }); diff --git a/apps/server/src/vcs/GitVcsDriver.ts b/apps/server/src/vcs/GitVcsDriver.ts index 4d63a447d63f..b5bd9aeb484a 100644 --- a/apps/server/src/vcs/GitVcsDriver.ts +++ b/apps/server/src/vcs/GitVcsDriver.ts @@ -48,6 +48,12 @@ export interface ExecuteGitInput { readonly timeoutMs?: number | null; readonly maxOutputBytes?: number; readonly appendTruncationMarker?: boolean; + /** + * With `appendTruncationMarker`, keep invoking the line callbacks after the + * buffered copy is full. For long-running commands whose output is only + * consumed through `progress`. + */ + readonly keepLineCallbacksAfterTruncation?: boolean; readonly progress?: ExecuteGitProgress; } diff --git a/apps/server/src/vcs/GitVcsDriverCore.test.ts b/apps/server/src/vcs/GitVcsDriverCore.test.ts index bc8a12700997..f1da4a6bda33 100644 --- a/apps/server/src/vcs/GitVcsDriverCore.test.ts +++ b/apps/server/src/vcs/GitVcsDriverCore.test.ts @@ -770,6 +770,38 @@ it.layer(TestLayer)("GitVcsDriver core integration", (it) => { }), ); + it.effect("keeps line callbacks flowing past the output cap when asked", () => + Effect.gen(function* () { + const cwd = yield* makeTmpDir(); + const driver = yield* GitVcsDriver.GitVcsDriver; + // 4 KiB of multi-byte lines, well past a 512-byte cap; the last line + // is the one a failure surface would need. + const lines: Array = []; + const result = yield* driver.execute({ + operation: "GitVcsDriver.test.callbacksPastCap", + cwd, + args: [ + "-c", + 'alias.spew=!for i in $(seq 1 128); do printf "é%03d\\n" $i >&2; done; echo fatal: last line >&2', + "spew", + ], + maxOutputBytes: 512, + appendTruncationMarker: true, + keepLineCallbacksAfterTruncation: true, + progress: { onStderrLine: (line) => Effect.sync(() => void lines.push(line)) }, + }); + + assert.isTrue(result.stderrTruncated); + assert.isAtMost(result.stderr.length, 600); + assert.equal(lines.length, 129); + assert.equal(lines[0], "é001"); + assert.equal(lines[127], "é128"); + assert.equal(lines.at(-1), "fatal: last line"); + // No replacement characters: the cap landing inside "é" is invisible to callbacks. + assert.isFalse(lines.some((line) => line.includes("\uFFFD"))); + }), + ); + it.effect("recovers a structurally identified missing cwd as a non-repository", () => Effect.gen(function* () { const parent = yield* makeTmpDir(); diff --git a/apps/server/src/vcs/GitVcsDriverCore.ts b/apps/server/src/vcs/GitVcsDriverCore.ts index 86a3e2ebd812..28ae6c5288ae 100644 --- a/apps/server/src/vcs/GitVcsDriverCore.ts +++ b/apps/server/src/vcs/GitVcsDriverCore.ts @@ -666,12 +666,19 @@ const collectOutput = Effect.fnUntraced(function* ( maxOutputBytes: number, appendTruncationMarker: boolean, onLine: ((line: string) => Effect.Effect) | undefined, + keepLineCallbacksAfterTruncation = false, ): Effect.fn.Return<{ readonly text: string; readonly truncated: boolean }, GitCommandError> { const decoder = new TextDecoder(); + // With callbacks continuing past the cap, lines are decoded by their own + // decoder from the first byte so no character is ever split at the cap. + const lineDecoder = keepLineCallbacksAfterTruncation && onLine ? new TextDecoder() : null; let bytes = 0; let text = ""; let lineBuffer = ""; let truncated = false; + // A separator-free stream past the cap must not grow the line buffer + // without bound; a line longer than this is not one the callbacks want. + const maxPendingLineBytes = 64 * 1024; // Git redraws progress with a bare `\r` between updates and only ends the // line once the step is done, so `\r` has to count as a line break here. @@ -697,6 +704,11 @@ const collectOutput = Effect.fnUntraced(function* ( const processChunk = Effect.fnUntraced(function* (chunk: Uint8Array) { if (appendTruncationMarker && truncated) { + if (lineDecoder) { + lineBuffer += lineDecoder.decode(chunk, { stream: true }); + yield* emitCompleteLines(false); + if (lineBuffer.length > maxPendingLineBytes) lineBuffer = ""; + } return; } const nextBytes = bytes + chunk.byteLength; @@ -717,7 +729,7 @@ const collectOutput = Effect.fnUntraced(function* ( const decoded = decoder.decode(chunkToDecode, { stream: !truncated }); text += decoded; - lineBuffer += decoded; + lineBuffer += lineDecoder ? lineDecoder.decode(chunk, { stream: true }) : decoded; yield* emitCompleteLines(false); }); @@ -735,6 +747,7 @@ const collectOutput = Effect.fnUntraced(function* ( const remainder = truncated ? "" : decoder.decode(); text += remainder; lineBuffer += remainder; + if (lineDecoder) lineBuffer += lineDecoder.decode(); yield* emitCompleteLines(true); return { text, @@ -802,6 +815,7 @@ export const makeGitVcsDriverCore = Effect.fn("makeGitVcsDriverCore")(function* maxOutputBytes, appendTruncationMarker, input.progress?.onStdoutLine, + input.keepLineCallbacksAfterTruncation, ), collectOutput( commandInput, @@ -809,6 +823,7 @@ export const makeGitVcsDriverCore = Effect.fn("makeGitVcsDriverCore")(function* maxOutputBytes, appendTruncationMarker, input.progress?.onStderrLine, + input.keepLineCallbacksAfterTruncation, ), child.exitCode.pipe( Effect.mapError( diff --git a/apps/server/src/ws.ts b/apps/server/src/ws.ts index 6a16640f7034..6e5d9c02db39 100644 --- a/apps/server/src/ws.ts +++ b/apps/server/src/ws.ts @@ -8,6 +8,7 @@ import * as DateTime from "effect/DateTime"; import * as Duration from "effect/Duration"; import * as Effect from "effect/Effect"; import * as Fiber from "effect/Fiber"; +import * as FileSystem from "effect/FileSystem"; import * as Layer from "effect/Layer"; import * as Option from "effect/Option"; import * as Path from "effect/Path"; @@ -133,6 +134,8 @@ import * as GitWorkflowService from "./git/GitWorkflowService.ts"; import { linkCreatedPullRequest } from "./git/linkCreatedPullRequest.ts"; import * as ReviewService from "./review/ReviewService.ts"; import * as ProjectSetupScriptRunner from "./project/ProjectSetupScriptRunner.ts"; +import * as ProjectCloneTracker from "./project/ProjectCloneTracker.ts"; +import * as RepositoryIdentityResolver from "./project/RepositoryIdentityResolver.ts"; import * as WorktreeSetupTracker from "./project/WorktreeSetupTracker.ts"; import * as AgentSessionScanner from "./project/AgentSessionScanner.ts"; import { importRecentAgentThreads } from "./project/AgentSessionImporter.ts"; @@ -605,6 +608,17 @@ const makeWsRpcLayer = ( }); const projectSetupScriptRunner = yield* ProjectSetupScriptRunner.ProjectSetupScriptRunner; const worktreeSetupTracker = yield* WorktreeSetupTracker.WorktreeSetupTracker; + const projectCloneTracker = yield* ProjectCloneTracker.ProjectCloneTracker; + const repositoryIdentityResolver = + yield* RepositoryIdentityResolver.RepositoryIdentityResolver; + // Clone hooks run on the tracker's fiber, outside any RPC, so the + // normalizer's services are captured here rather than inherited. + const normalizerContext = yield* Effect.context< + | FileSystem.FileSystem + | Path.Path + | ServerConfig.ServerConfig + | WorkspacePaths.WorkspacePaths + >(); const agentSessionScanner = yield* AgentSessionScanner.AgentSessionScanner; const serverEnvironment = yield* ServerEnvironment.ServerEnvironment; const backgroundPolicy = yield* BackgroundPolicy.BackgroundPolicy; @@ -1615,6 +1629,7 @@ const makeWsRpcLayer = ( observeRpcEffect( ORCHESTRATION_WS_METHODS.dispatchCommand, Effect.gen(function* () { + yield* ProjectCloneTracker.rejectCommandsDuringClone(projectCloneTracker, command); const normalizedCommand = yield* normalizeDispatchCommand(command); // Archive removes the thread from the client, so this transport // closes its session and terminals after the command lands. @@ -1646,6 +1661,10 @@ const makeWsRpcLayer = ( Effect.tapError(() => cleanupFailedUploadedAttachments(command, normalizedCommand)), ); yield* recordClientCommandAnalytics(normalizedCommand); + yield* ProjectCloneTracker.discardCloneForDeletedProject( + projectCloneTracker, + normalizedCommand, + ); if (archiveCommand) { if (shouldStopSessionAfterCommand) { yield* Effect.gen(function* () { @@ -2657,6 +2676,65 @@ const makeWsRpcLayer = ( "rpc.aggregate": "source-control", }, ), + [WS_METHODS.projectCloneStart]: (input) => + observeRpcEffect( + WS_METHODS.projectCloneStart, + projectCloneTracker.start(input, { + createProject: (project) => + Effect.gen(function* () { + const normalizedCommand = yield* normalizeDispatchCommand({ + type: "project.create", + commandId: yield* serverCommandId("project-clone-create"), + projectId: project.projectId, + title: project.title, + workspaceRoot: project.workspaceRoot, + createWorkspaceRootIfMissing: true, + createdAt: project.createdAt, + }); + yield* dispatchNormalizedCommand(normalizedCommand); + yield* recordClientCommandAnalytics(normalizedCommand); + }).pipe(Effect.provideContext(normalizerContext)), + onCloned: (project) => + // The project was created against an empty directory, so its + // cached identity is "not a repository" until this refresh. + // Re-emitting the project shell carries the new identity to + // every client without a round trip. + repositoryIdentityResolver.resolve(project.workspaceRoot, { refresh: true }).pipe( + Effect.andThen( + Effect.gen(function* () { + const command = yield* normalizeDispatchCommand({ + type: "project.meta.update", + commandId: yield* serverCommandId("project-clone-done"), + projectId: project.projectId, + }); + yield* dispatchNormalizedCommand(command); + }), + ), + Effect.andThen(refreshGitStatus(project.workspaceRoot)), + Effect.ignoreCause({ log: true }), + Effect.provideContext(normalizerContext), + ), + }), + { "rpc.aggregate": "source-control" }, + ), + [WS_METHODS.projectCloneCancel]: (input) => + observeRpcEffect( + WS_METHODS.projectCloneCancel, + projectCloneTracker + .cancel(input.projectId) + .pipe(Effect.map((applied) => ({ applied }))), + { "rpc.aggregate": "source-control" }, + ), + [WS_METHODS.projectCloneRetry]: (input) => + observeRpcEffect( + WS_METHODS.projectCloneRetry, + projectCloneTracker.retry(input.projectId).pipe(Effect.map((applied) => ({ applied }))), + { "rpc.aggregate": "source-control" }, + ), + [WS_METHODS.subscribeProjectClones]: () => + observeRpcStream(WS_METHODS.subscribeProjectClones, projectCloneTracker.stream, { + "rpc.aggregate": "source-control", + }), [WS_METHODS.sourceControlPublishRepository]: (input) => observeRpcEffect( WS_METHODS.sourceControlPublishRepository, diff --git a/apps/web/src/components/ChatView.tsx b/apps/web/src/components/ChatView.tsx index 0b6262387a96..303427ebd2e7 100644 --- a/apps/web/src/components/ChatView.tsx +++ b/apps/web/src/components/ChatView.tsx @@ -226,6 +226,7 @@ import { AlarmClockIcon, CheckCircle2Icon, ChevronDownIcon, + DownloadIcon, GitBranchIcon, Minimize2Icon, PaperclipIcon, @@ -259,6 +260,7 @@ import { import { useNowMinute } from "../hooks/useNowMinute"; import { usePanelAnimationSettings, usePanelPresence } from "../panelAnimations"; import { useNewThreadHandler } from "../hooks/useHandleNewThread"; +import { useRemoveClonedProject } from "../hooks/useRemoveClonedProject"; import { useOpenPanelPullRequestUrl } from "../hooks/useOpenPanelPullRequestUrl"; import { useThreadActions } from "../hooks/useThreadActions"; import { resolveAppModelSelectionForInstance } from "../modelSelection"; @@ -328,6 +330,9 @@ import { } from "@t3tools/client-runtime/state/threads"; import { resolveProviderSkillsForCwd } from "@t3tools/client-runtime/providerSkills"; import { vcsEnvironment } from "../state/vcs"; +import { sourceControlEnvironment } from "../state/sourceControl"; +import { useProjectClone } from "../state/projectClones"; +import { projectCloneDisplayName, projectCloneProgressSummary } from "@t3tools/contracts"; import { useEnvironments, usePrimaryEnvironment } from "../state/environments"; import { useProject, @@ -2094,6 +2099,113 @@ export default function ChatView(props: ChatViewProps) { () => (activeProject ? resolveProjectScripts(settings, activeProject) : []), [activeProject, settings], ); + // A project added by cloning exists before its files do. The draft stays + // editable throughout; only sending waits for the clone, and a failed + // clone offers its retry right where the user is looking. + const activeProjectClone = useProjectClone(activeProjectRef); + const cancelProjectClone = useAtomCommand(sourceControlEnvironment.cancelProjectClone, { + reportFailure: false, + }); + const retryProjectClone = useAtomCommand(sourceControlEnvironment.retryProjectClone, { + reportFailure: false, + }); + const removeClonedProject = useRemoveClonedProject(); + // The banner mirrors the server's clone state, so a request that never got + // there needs its own feedback. + const runProjectCloneAction = useCallback( + async ( + title: string, + action: () => Promise>, + ): Promise => { + const result = await action(); + if (result._tag === "Failure" && !isAtomCommandInterrupted(result)) { + const error = squashAtomCommandFailure(result); + toastManager.add( + stackedThreadToast({ + type: "error", + title, + description: error instanceof Error ? error.message : "An error occurred.", + }), + ); + } + }, + [], + ); + const projectCloneSendBlockReason = + activeProjectClone === null + ? null + : activeProjectClone.phase === "running" + ? "Cloning repository" + : activeProjectClone.phase === "done" + ? null + : "Repository not cloned"; + const projectCloneBannerItem = useMemo(() => { + if (!activeProjectClone || !activeProjectRef || activeProjectClone.phase === "done") { + return null; + } + const name = projectCloneDisplayName(activeProjectClone); + const { environmentId, projectId } = activeProjectRef; + if (activeProjectClone.phase === "running") { + return { + id: `project-clone:${projectId}`, + variant: "info", + priority: "activity", + icon: , + title: `Cloning ${name}`, + description: projectCloneProgressSummary(activeProjectClone), + actions: ( + + ), + }; + } + const cancelled = activeProjectClone.phase === "cancelled"; + return { + id: `project-clone:${projectId}`, + variant: cancelled ? "warning" : "error", + icon: , + title: cancelled ? `Cancelled cloning ${name}` : `Failed to clone ${name}`, + description: cancelled ? "Retry to bring in the repository." : activeProjectClone.error, + actions: ( + <> + + + + ), + }; + }, [ + activeProjectClone, + activeProjectRef, + cancelProjectClone, + removeClonedProject, + retryProjectClone, + runProjectCloneAction, + ]); const activeProjectDefaultModelSelection = activeProjectSettings.settings.defaultModelSelection; const handleNewThreadInActiveProject = useCallback(() => { startNewThreadForProject(activeProjectRef, handleNewThread); @@ -6248,10 +6360,12 @@ export default function ChatView(props: ChatViewProps) { const parkedThreadItems = parkedThreadBannerItem === null ? [] : [parkedThreadBannerItem]; // The user asked for this one, so it leads the notice tier instead of trailing it. const usageLimitsItems = usageLimitsBanner === null ? [] : [usageLimitsBanner]; + const projectCloneItems = projectCloneBannerItem === null ? [] : [projectCloneBannerItem]; if (!localCheckoutBranchMismatch || !showBranchMismatchBanner || !activeBranchMismatchKey) { return [ ...feedbackBannerItems, ...usageLimitsItems, + ...projectCloneItems, ...systemComposerBannerItems, ...backgroundLivenessItems, ...resumeCompactionItems, @@ -6262,6 +6376,7 @@ export default function ChatView(props: ChatViewProps) { return [ ...feedbackBannerItems, ...usageLimitsItems, + ...projectCloneItems, ...systemComposerBannerItems, ...backgroundLivenessItems, ...resumeCompactionItems, @@ -6314,6 +6429,7 @@ export default function ChatView(props: ChatViewProps) { isRestoringThreadBranch, localCheckoutBranchMismatch, parkedThreadBannerItem, + projectCloneBannerItem, resumeCompactionBannerItem, showBranchMismatchBanner, systemComposerBannerItems, @@ -9162,7 +9278,7 @@ export default function ChatView(props: ChatViewProps) { ? "Sending feedback" : threadDetailLoading ? "Messages loading" - : null + : projectCloneSendBlockReason } isPreparingWorktree={isPreparingWorktree} bannerItems={composerBannerItems} diff --git a/apps/web/src/components/CommandPalette.tsx b/apps/web/src/components/CommandPalette.tsx index f33a9cce63b0..a8d1e57e5172 100644 --- a/apps/web/src/components/CommandPalette.tsx +++ b/apps/web/src/components/CommandPalette.tsx @@ -84,7 +84,7 @@ import { sourceControlEnvironment } from "../state/sourceControl"; import { useAtomCommand } from "../state/use-atom-command"; import { useAtomQueryRunner } from "../state/use-atom-query-runner"; import { useEnvironments, usePrimaryEnvironmentId } from "../state/environments"; -import { useProjects, useServerConfigs, useThreadShells } from "../state/entities"; +import { useProjects, useServerConfigs, useThreadShells, waitForProject } from "../state/entities"; import { useThreadSearch } from "../state/queries"; import { resolveThreadActionProjectRef, startNewThreadFromContext } from "../lib/chatThreadActions"; import { @@ -643,6 +643,9 @@ function OpenCommandPaletteDialog(props: { const cloneRepository = useAtomCommand(sourceControlEnvironment.cloneRepository, { reportFailure: false, }); + const startProjectClone = useAtomCommand(sourceControlEnvironment.startProjectClone, { + reportFailure: false, + }); const { environments } = useEnvironments(); const desktopLocalBootstraps = useDesktopLocalBootstraps(); const primaryEnvironmentId = usePrimaryEnvironmentId(); @@ -2198,28 +2201,80 @@ function OpenCommandPaletteDialog(props: { return; } + // Older servers only offer the blocking clone: the palette has to wait + // for git so it can add the project afterwards. + if (browseEnvironment?.serverConfig?.environment.capabilities.projectCloneTracking !== true) { + setIsRemoteProjectCloning(true); + const cloneResult = await cloneRepository({ + environmentId: addProjectCloneFlow.environmentId, + input: { + remoteUrl: addProjectCloneFlow.remoteUrl, + destinationPath, + }, + }); + setIsRemoteProjectCloning(false); + if (cloneResult._tag === "Failure") { + if (!isAtomCommandInterrupted(cloneResult)) { + toastManager.add( + stackedThreadToast({ + type: "error", + title: "Clone failed", + description: errorMessage(squashAtomCommandFailure(cloneResult)), + }), + ); + } + return; + } + await handleAddProject(cloneResult.value.cwd); + return; + } + + // The server creates the project and clones in the background; progress + // shows in a toast and in the draft's composer banner, so the palette + // closes as soon as the clone is under way. Only problems found before + // git runs (bad destination, unknown repository) come back here. + const projectId = newProjectId(); setIsRemoteProjectCloning(true); - const cloneResult = await cloneRepository({ + const startResult = await startProjectClone({ environmentId: addProjectCloneFlow.environmentId, input: { + projectId, + title: inferProjectTitleFromPath(destinationPath), + createdAt: new Date().toISOString(), remoteUrl: addProjectCloneFlow.remoteUrl, destinationPath, }, }); setIsRemoteProjectCloning(false); - if (cloneResult._tag === "Failure") { - if (!isAtomCommandInterrupted(cloneResult)) { + if (startResult._tag === "Failure") { + if (!isAtomCommandInterrupted(startResult)) { toastManager.add( stackedThreadToast({ type: "error", title: "Clone failed", - description: errorMessage(squashAtomCommandFailure(cloneResult)), + description: errorMessage(squashAtomCommandFailure(startResult)), }), ); } return; } - await handleAddProject(cloneResult.value.cwd); + setOpen(false); + const projectRef = scopeProjectRef(addProjectCloneFlow.environmentId, projectId); + // The create event usually lands before this call returns; give the shell + // stream a moment so the draft opens with its project resolved instead of + // flashing the project picker. + await waitForProject(projectRef, 3_000).catch(() => null); + const navigationResult = await settlePromise(() => handleNewThread(projectRef)); + if (navigationResult._tag === "Failure") { + const error = squashAtomCommandFailure(navigationResult); + toastManager.add( + stackedThreadToast({ + type: "error", + title: "Failed to open project", + description: error instanceof Error ? error.message : "An error occurred.", + }), + ); + } } const browseTo = useCallback( diff --git a/apps/web/src/components/ProjectCloneToastCoordinator.tsx b/apps/web/src/components/ProjectCloneToastCoordinator.tsx new file mode 100644 index 000000000000..2c4a477a2f2e --- /dev/null +++ b/apps/web/src/components/ProjectCloneToastCoordinator.tsx @@ -0,0 +1,240 @@ +import { useParams } from "@tanstack/react-router"; +import { scopeProjectRef } from "@t3tools/client-runtime/environment"; +import { + type AtomCommandResult, + isAtomCommandInterrupted, + squashAtomCommandFailure, +} from "@t3tools/client-runtime/state/runtime"; +import { + projectCloneDisplayName, + projectCloneProgressSummary, + type EnvironmentId, + type ProjectCloneSnapshot, + type ProjectId, +} from "@t3tools/contracts"; +import { useCallback, useEffect, useRef } from "react"; + +import { useNewThreadHandler } from "../hooks/useHandleNewThread"; +import { useRemoveClonedProject } from "../hooks/useRemoveClonedProject"; +import { useEnvironments } from "../state/environments"; +import { useEnvironmentProjectClones } from "../state/projectClones"; +import { sourceControlEnvironment } from "../state/sourceControl"; +import { useAtomCommand } from "../state/use-atom-command"; +import { type DraftId, useComposerDraftStore } from "../composerDraftStore"; +import { toastManager } from "./ui/toast"; +import { stackedThreadToast } from "./ui/toastHelpers"; + +/** + * One toast per clone in flight, on every environment. The palette that + * started a clone closes right away, so this is where its progress lives: + * the toast updates in place as git reports stages, then settles into a + * success or failure state with the matching action. + */ +export function ProjectCloneToastCoordinator() { + const { environments } = useEnvironments(); + return environments.map((environment) => ( + + )); +} + +interface TrackedToast { + readonly toastId: ReturnType; + /** The last snapshot rendered, so an identical redraw does not touch the toast. */ + readonly renderedKey: string; + readonly phase: ProjectCloneSnapshot["phase"]; +} + +function renderKey(clone: ProjectCloneSnapshot): string { + return `${clone.phase}:${clone.stage}:${clone.percent ?? ""}:${clone.detail ?? ""}:${clone.error ?? ""}`; +} + +function EnvironmentCloneToasts({ environmentId }: { environmentId: EnvironmentId }) { + const clones = useEnvironmentProjectClones(environmentId); + const handleNewThread = useNewThreadHandler(); + const { draftId: routeDraftId } = useParams({ strict: false }); + const cancelClone = useAtomCommand(sourceControlEnvironment.cancelProjectClone, { + reportFailure: false, + }); + const retryClone = useAtomCommand(sourceControlEnvironment.retryProjectClone, { + reportFailure: false, + }); + // The toast mirrors the server's clone state, so a request that never got + // there needs its own feedback. + const runCloneAction = useCallback( + async (title: string, action: () => Promise>) => { + const result = await action(); + if (result._tag === "Failure" && !isAtomCommandInterrupted(result)) { + const error = squashAtomCommandFailure(result); + toastManager.add( + stackedThreadToast({ + type: "error", + title, + description: error instanceof Error ? error.message : "An error occurred.", + }), + ); + } + }, + [], + ); + const removeClonedProject = useRemoveClonedProject(); + const toasts = useRef(new Map()); + + // Whether the user is already looking at this project's draft: the composer + // banner shows the same progress and actions there, so the toast steps + // aside and comes back if they navigate away mid-clone. + const isViewingProjectDraft = useCallback( + (projectId: ProjectId) => { + if (!routeDraftId) return false; + const draft = useComposerDraftStore.getState().getDraftSession(routeDraftId as DraftId); + return draft?.environmentId === environmentId && draft.projectId === projectId; + }, + [environmentId, routeDraftId], + ); + + const openProject = useCallback( + (projectId: ProjectId) => { + void handleNewThread(scopeProjectRef(environmentId, projectId)); + }, + [environmentId, handleNewThread], + ); + + useEffect(() => { + const seen = new Set(); + for (const clone of clones) { + seen.add(clone.projectId); + const key = renderKey(clone); + const tracked = toasts.current.get(clone.projectId); + const name = projectCloneDisplayName(clone); + // Handlers run later than this pass, so they look the toast up then. + const closeToast = () => { + const current = toasts.current.get(clone.projectId); + if (!current) return; + toastManager.close(current.toastId); + toasts.current.delete(clone.projectId); + }; + if (isViewingProjectDraft(clone.projectId)) { + closeToast(); + continue; + } + if (tracked?.renderedKey === key) continue; + + if (clone.phase === "running") { + const options = stackedThreadToast({ + type: "loading", + title: `Cloning ${name}`, + description: projectCloneProgressSummary(clone), + timeout: 0, + actionProps: { + children: "Cancel", + onClick: () => { + void runCloneAction("Failed to cancel clone", () => + cancelClone({ environmentId, input: { projectId: clone.projectId } }), + ); + }, + }, + data: { hideCopyButton: true }, + }); + if (tracked) { + toastManager.update(tracked.toastId, options); + toasts.current.set(clone.projectId, { ...tracked, renderedKey: key, phase: "running" }); + } else { + const toastId = toastManager.add(options); + toasts.current.set(clone.projectId, { toastId, renderedKey: key, phase: "running" }); + } + continue; + } + + if (clone.phase === "done") { + const options = stackedThreadToast({ + type: "success", + title: `Cloned ${name}`, + description: clone.destinationPath, + timeout: 8_000, + actionProps: { + children: "Open project", + onClick: () => { + closeToast(); + openProject(clone.projectId); + }, + }, + data: { hideCopyButton: true }, + }); + if (tracked) { + toastManager.update(tracked.toastId, options); + toasts.current.set(clone.projectId, { ...tracked, renderedKey: key, phase: "done" }); + } else { + const toastId = toastManager.add(options); + toasts.current.set(clone.projectId, { toastId, renderedKey: key, phase: "done" }); + } + continue; + } + + // Failed or cancelled: the project stays, pointing at an empty folder. + // Retry from here; the draft's composer banner offers the same. + const cancelled = clone.phase === "cancelled"; + const options = stackedThreadToast({ + type: cancelled ? "info" : "error", + title: cancelled ? `Cancelled cloning ${name}` : `Failed to clone ${name}`, + description: cancelled ? clone.destinationPath : (clone.error ?? "The clone failed."), + timeout: 0, + actionProps: { + children: "Retry", + onClick: () => { + void runCloneAction("Failed to retry clone", () => + retryClone({ environmentId, input: { projectId: clone.projectId } }), + ); + }, + }, + data: { + ...(cancelled ? { hideCopyButton: true } : {}), + secondaryActionProps: { + children: "Remove project", + onClick: () => { + // The server drops the clone with the project, which closes + // this toast; a failed removal leaves it (and Retry) in place. + void removeClonedProject({ environmentId, projectId: clone.projectId }); + }, + }, + }, + }); + if (tracked) { + toastManager.update(tracked.toastId, options); + toasts.current.set(clone.projectId, { ...tracked, renderedKey: key, phase: clone.phase }); + } else { + const toastId = toastManager.add(options); + toasts.current.set(clone.projectId, { toastId, renderedKey: key, phase: clone.phase }); + } + } + + // A clone the server stopped tracking (done and expired, or its project + // was removed) takes its toast with it, unless it already settled into a + // timed success toast that dismisses itself. + for (const [projectId, tracked] of toasts.current) { + if (seen.has(projectId)) continue; + if (tracked.phase !== "done") toastManager.close(tracked.toastId); + toasts.current.delete(projectId); + } + }, [ + cancelClone, + clones, + environmentId, + isViewingProjectDraft, + openProject, + removeClonedProject, + retryClone, + runCloneAction, + ]); + + useEffect( + () => () => { + for (const tracked of toasts.current.values()) toastManager.close(tracked.toastId); + toasts.current.clear(); + }, + [], + ); + + return null; +} diff --git a/apps/web/src/hooks/useRemoveClonedProject.ts b/apps/web/src/hooks/useRemoveClonedProject.ts new file mode 100644 index 000000000000..72c5b91cde32 --- /dev/null +++ b/apps/web/src/hooks/useRemoveClonedProject.ts @@ -0,0 +1,62 @@ +import { useRouter } from "@tanstack/react-router"; +import { squashAtomCommandFailure } from "@t3tools/client-runtime/state/runtime"; +import type { ScopedProjectRef } from "@t3tools/contracts"; +import { useCallback } from "react"; + +import { useComposerDraftStore } from "../composerDraftStore"; +import { resolveThreadRouteTarget } from "../threadRoutes"; +import { releaseProjectDraftUploads } from "../lib/composerDraftUploads"; +import { projectEnvironment } from "../state/projects"; +import { useAtomCommand } from "../state/use-atom-command"; +import { stackedThreadToast, toastManager } from "../components/ui/toast"; + +/** + * Removes a project whose clone never landed. The server clears the empty + * folder along with the clone, so there is nothing to confirm: no threads + * exist yet and the draft is the only thing lost, which the user is looking + * at when they click. + */ +export function useRemoveClonedProject() { + const router = useRouter(); + const deleteProject = useAtomCommand(projectEnvironment.delete, { reportFailure: false }); + + return useCallback( + async (projectRef: ScopedProjectRef) => { + const draftStore = useComposerDraftStore.getState(); + const result = await deleteProject({ + environmentId: projectRef.environmentId, + // Not forced: a project whose clone never landed has no threads, and + // if one appeared in the meantime the server refuses rather than + // silently deleting it. + input: { projectId: projectRef.projectId }, + }); + if (result._tag === "Failure") { + const error = squashAtomCommandFailure(result); + toastManager.add( + stackedThreadToast({ + type: "error", + title: "Failed to remove project", + description: error instanceof Error ? error.message : "An error occurred.", + }), + ); + return false; + } + // Read the route after the await: the user may have moved on while the + // delete was in flight, and only a draft of this project needs to go. + const routeParams = router.state.matches[router.state.matches.length - 1]?.params ?? {}; + const routeTarget = resolveThreadRouteTarget(routeParams); + const viewingDraft = + routeTarget?.kind === "draft" ? draftStore.getDraftSession(routeTarget.draftId) : null; + const viewingThisProject = + viewingDraft?.environmentId === projectRef.environmentId && + viewingDraft.projectId === projectRef.projectId; + releaseProjectDraftUploads(projectRef); + const projectDraft = draftStore.getDraftThreadByProjectRef(projectRef); + if (projectDraft) draftStore.clearDraftThread(projectDraft.draftId); + draftStore.clearProjectDraftThreadId(projectRef); + if (viewingThisProject) void router.navigate({ to: "/", replace: true }); + return true; + }, + [deleteProject, router], + ); +} diff --git a/apps/web/src/routes/__root.tsx b/apps/web/src/routes/__root.tsx index 89a1ff23cc83..4db30275ee6b 100644 --- a/apps/web/src/routes/__root.tsx +++ b/apps/web/src/routes/__root.tsx @@ -27,6 +27,7 @@ import { SnapShotCoordinator } from "../components/desktop/SnapShotCoordinator"; import { DesktopAppActivationCoordinator } from "../components/desktop/DesktopAppActivationCoordinator"; import { ProviderUpdateLaunchNotification } from "../components/ProviderUpdateLaunchNotification"; import { ThreadNotificationCoordinator } from "../components/ThreadNotificationCoordinator"; +import { ProjectCloneToastCoordinator } from "../components/ProjectCloneToastCoordinator"; import { SlowRpcRequestToastCoordinator } from "../components/SlowRpcRequestToastCoordinator"; import { ThemeEditorHost } from "../components/settings/ThemeEditorHost"; import { useCopyToClipboard } from "../hooks/useCopyToClipboard"; @@ -222,6 +223,7 @@ function RootRouteView() { + {primaryEnvironmentAuthenticated ? ( diff --git a/apps/web/src/state/projectClones.ts b/apps/web/src/state/projectClones.ts new file mode 100644 index 000000000000..b5645167c1cc --- /dev/null +++ b/apps/web/src/state/projectClones.ts @@ -0,0 +1,52 @@ +import { useAtomValue } from "@effect/atom-react"; +import { parseScopedProjectKey, scopedProjectKey } from "@t3tools/client-runtime/environment"; +import type { EnvironmentId, ProjectCloneSnapshot, ScopedProjectRef } from "@t3tools/contracts"; +import * as Option from "effect/Option"; +import { AsyncResult, Atom } from "effect/unstable/reactivity"; + +import { environmentServerConfigsAtom } from "./server"; +import { sourceControlEnvironment } from "./sourceControl"; + +const EMPTY_CLONES: ReadonlyArray = []; +const EMPTY_CLONE_ATOM = Atom.make(null).pipe( + Atom.withLabel("web-project-clone:empty"), +); + +/** + * Latest clone list an environment has streamed; empty until the subscription + * delivers, and never subscribed on servers that predate clone tracking. + */ +const environmentProjectClonesAtom = Atom.family((environmentId: EnvironmentId) => + Atom.make((get): ReadonlyArray => { + const supported = + get(environmentServerConfigsAtom).get(environmentId)?.environment.capabilities + .projectCloneTracking === true; + if (!supported) return EMPTY_CLONES; + const result = get(sourceControlEnvironment.projectClones({ environmentId, input: {} })); + return Option.getOrElse(AsyncResult.value(result), () => EMPTY_CLONES); + }).pipe(Atom.withLabel(`web-project-clones:${environmentId}`)), +); + +const projectCloneAtom = Atom.family((key: string) => { + const ref = parseScopedProjectKey(key); + return Atom.make((get): ProjectCloneSnapshot | null => { + if (ref === null) return null; + const clones = get(environmentProjectClonesAtom(ref.environmentId)); + return clones.find((clone) => clone.projectId === ref.projectId) ?? null; + }).pipe(Atom.withLabel(`web-project-clone:${key}`)); +}); + +/** + * The tracked clone for a project, or null once it finished (or never + * existed). Subscribing here opens the environment's clone stream, which is + * cheap: the server sends an empty list and stays quiet until a clone starts. + */ +export function useProjectClone(ref: ScopedProjectRef | null): ProjectCloneSnapshot | null { + return useAtomValue(ref === null ? EMPTY_CLONE_ATOM : projectCloneAtom(scopedProjectKey(ref))); +} + +export function useEnvironmentProjectClones( + environmentId: EnvironmentId, +): ReadonlyArray { + return useAtomValue(environmentProjectClonesAtom(environmentId)); +} diff --git a/docs/user/source-control.md b/docs/user/source-control.md index 64cdc19c2812..d8a06fc34b64 100644 --- a/docs/user/source-control.md +++ b/docs/user/source-control.md @@ -76,7 +76,10 @@ az login ## Clone or publish a project Use **Add Project** in the command palette (`Cmd/Ctrl+K`) to clone a repository. Choose a hosting -provider or paste a Git URL, then choose where to save it. +provider or paste a Git URL, then choose where to save it. The project opens right away while the +clone runs in the background: you can write your first prompt, and sending waits until the files +are in place. A toast tracks progress and lets you cancel; if the clone fails, retry it from the +toast or from the banner above the composer. For a local Git repository without a remote, **Publish Repository** creates a hosted repository, adds it as `origin`, and pushes your commits. If there are no commits yet, it creates the remote; diff --git a/packages/client-runtime/src/rpc/client.ts b/packages/client-runtime/src/rpc/client.ts index af140ef2fcde..e80a2f0f4b12 100644 --- a/packages/client-runtime/src/rpc/client.ts +++ b/packages/client-runtime/src/rpc/client.ts @@ -57,6 +57,7 @@ export type EnvironmentSubscriptionRpcTag = | typeof WS_METHODS.previewAutomationConnect | typeof WS_METHODS.subscribeVcsStatus | typeof WS_METHODS.subscribeWorktreeSetup + | typeof WS_METHODS.subscribeProjectClones | typeof WS_METHODS.terminalAttach; export type EnvironmentStreamCommandRpcTag = diff --git a/packages/client-runtime/src/state/sourceControl.ts b/packages/client-runtime/src/state/sourceControl.ts index c1598b49eaeb..39c7a9544549 100644 --- a/packages/client-runtime/src/state/sourceControl.ts +++ b/packages/client-runtime/src/state/sourceControl.ts @@ -5,6 +5,7 @@ import { createAtomCommandScheduler, createEnvironmentRpcCommand, createEnvironmentRpcQueryAtomFamily, + createEnvironmentRpcSubscriptionAtomFamily, } from "./runtime.ts"; import type { EnvironmentRegistry } from "../connection/registry.ts"; import { EnvironmentCacheStore } from "../platform/persistence.ts"; @@ -33,6 +34,43 @@ export function createSourceControlEnvironmentAtoms( key: ({ environmentId }) => environmentId, }, }), + // Clone-backed project creation. The RPC returns once the project exists + // and the clone runs in the background; `projectClones` carries progress. + startProjectClone: createEnvironmentRpcCommand(runtime, { + label: "environment-data:source-control:project-clone-start", + tag: WS_METHODS.projectCloneStart, + scheduler: commandScheduler, + concurrency: { + mode: "serial", + key: ({ environmentId }) => environmentId, + }, + }), + // Cancel and retry share the start queue so a double click cannot race + // two actions against the same clone. + cancelProjectClone: createEnvironmentRpcCommand(runtime, { + label: "environment-data:source-control:project-clone-cancel", + tag: WS_METHODS.projectCloneCancel, + scheduler: commandScheduler, + concurrency: { + mode: "serial", + key: ({ environmentId }) => environmentId, + }, + }), + retryProjectClone: createEnvironmentRpcCommand(runtime, { + label: "environment-data:source-control:project-clone-retry", + tag: WS_METHODS.projectCloneRetry, + scheduler: commandScheduler, + concurrency: { + mode: "serial", + key: ({ environmentId }) => environmentId, + }, + }), + // Every clone the environment tracks. Empty until a clone starts; a + // finished clone drops out after a grace period, a failed one stays. + projectClones: createEnvironmentRpcSubscriptionAtomFamily(runtime, { + label: "environment-data:source-control:project-clones", + tag: WS_METHODS.subscribeProjectClones, + }), publishRepository: createEnvironmentRpcCommand(runtime, { label: "environment-data:source-control:publish-repository", tag: WS_METHODS.sourceControlPublishRepository, diff --git a/packages/contracts/src/environment.ts b/packages/contracts/src/environment.ts index 88e2661523a3..c8b8833ead86 100644 --- a/packages/contracts/src/environment.ts +++ b/packages/contracts/src/environment.ts @@ -153,6 +153,11 @@ export const ExecutionEnvironmentCapabilities = Schema.Struct({ this is false — no update would ever repaint it. Absent on older servers, which may still publish, so only an explicit false skips. */ agentActivityPublishing: Schema.optionalKey(Schema.Boolean), + /** Server runs repository clones for new projects in the background and + streams their progress (`projectClone.*`, `subscribeProjectClones`). + Absent on older servers, where clients must clone with the blocking + `sourceControl.cloneRepository` call instead. */ + projectCloneTracking: Schema.optionalKey(Schema.Boolean), /** Server detects `platform.machine` and persists the `environmentIcon` setting. Older servers drop the key on write, so clients show the picker inert rather than offering a choice that would never stick. */ diff --git a/packages/contracts/src/index.ts b/packages/contracts/src/index.ts index 007120dcba05..978a0459e69b 100644 --- a/packages/contracts/src/index.ts +++ b/packages/contracts/src/index.ts @@ -25,6 +25,7 @@ export * from "./settings.ts"; export * from "./git.ts"; export * from "./vcs.ts"; export * from "./sourceControl.ts"; +export * from "./projectClone.ts"; export * from "./pullRequest.ts"; export * from "./orchestration.ts"; export * from "./t3ProjectFile.ts"; diff --git a/packages/contracts/src/projectClone.ts b/packages/contracts/src/projectClone.ts new file mode 100644 index 000000000000..e624ec17a72f --- /dev/null +++ b/packages/contracts/src/projectClone.ts @@ -0,0 +1,122 @@ +import * as Schema from "effect/Schema"; + +import { IsoDateTime, NonNegativeInt, ProjectId, TrimmedNonEmptyString } from "./baseSchemas.ts"; +import { + SourceControlCloneProtocol, + SourceControlProviderKind, + SourceControlRepositoryInfo, +} from "./sourceControl.ts"; + +/** + * Live progress of a repository clone that backs a freshly added project. The + * server keeps this in memory only: a finished clone is dropped after a short + * grace period, a failed one stays until it is retried or the project is + * removed, and a server restart forgets in-flight clones (the project keeps + * its empty workspace root and can be removed like any other project). + */ +/** Producers clamp free text to these before publishing so encoding never fails. */ +export const PROJECT_CLONE_DETAIL_MAX_LENGTH = 200; +export const PROJECT_CLONE_ERROR_MAX_LENGTH = 1000; + +/** Follows git's own clone phases as they appear on stderr, in order. */ +export const ProjectCloneStage = Schema.Literals([ + "connecting", + "counting", + "receiving", + "resolving", + "checkout", +]); +export type ProjectCloneStage = typeof ProjectCloneStage.Type; + +export const ProjectClonePhase = Schema.Literals(["running", "done", "failed", "cancelled"]); +export type ProjectClonePhase = typeof ProjectClonePhase.Type; + +export const ProjectCloneSnapshot = Schema.Struct({ + projectId: ProjectId, + remoteUrl: TrimmedNonEmptyString, + destinationPath: TrimmedNonEmptyString, + repository: Schema.NullOr(SourceControlRepositoryInfo), + phase: ProjectClonePhase, + stage: ProjectCloneStage, + /** Percent of the current stage, parsed from git's progress lines. */ + percent: Schema.NullOr(Schema.Int.check(Schema.isBetween({ minimum: 0, maximum: 100 }))), + /** Trailing text from the progress line, typically transfer size and rate. */ + detail: Schema.NullOr(Schema.String.check(Schema.isMaxLength(PROJECT_CLONE_DETAIL_MAX_LENGTH))), + /** Human readable reason when phase is failed. */ + error: Schema.NullOr(Schema.String.check(Schema.isMaxLength(PROJECT_CLONE_ERROR_MAX_LENGTH))), + startedAt: IsoDateTime, + endedAt: Schema.NullOr(IsoDateTime), + sequence: NonNegativeInt, +}); +export type ProjectCloneSnapshot = typeof ProjectCloneSnapshot.Type; + +export const ProjectCloneSubscribeInput = Schema.Struct({}); +export type ProjectCloneSubscribeInput = typeof ProjectCloneSubscribeInput.Type; + +/** Every tracked clone on the environment. Sent first, then after every change. */ +export const ProjectCloneListEvent = Schema.Array(ProjectCloneSnapshot); +export type ProjectCloneListEvent = typeof ProjectCloneListEvent.Type; + +export const ProjectCloneStartInput = Schema.Struct({ + projectId: ProjectId, + title: TrimmedNonEmptyString, + createdAt: IsoDateTime, + provider: Schema.optional(SourceControlProviderKind), + repository: Schema.optional(TrimmedNonEmptyString), + remoteUrl: Schema.optional(TrimmedNonEmptyString), + destinationPath: TrimmedNonEmptyString, + protocol: Schema.optional(SourceControlCloneProtocol), +}); +export type ProjectCloneStartInput = typeof ProjectCloneStartInput.Type; + +export const ProjectCloneStartResult = Schema.Struct({ + projectId: ProjectId, + cwd: TrimmedNonEmptyString, + remoteUrl: TrimmedNonEmptyString, + repository: Schema.NullOr(SourceControlRepositoryInfo), +}); +export type ProjectCloneStartResult = typeof ProjectCloneStartResult.Type; + +export const ProjectCloneActionInput = Schema.Struct({ + projectId: ProjectId, +}); +export type ProjectCloneActionInput = typeof ProjectCloneActionInput.Type; + +export const ProjectCloneActionResult = Schema.Struct({ + applied: Schema.Boolean, +}); +export type ProjectCloneActionResult = typeof ProjectCloneActionResult.Type; + +function projectCloneStageLabel(stage: ProjectCloneStage): string { + switch (stage) { + case "connecting": + return "Connecting"; + case "counting": + return "Counting objects"; + case "receiving": + return "Receiving objects"; + case "resolving": + return "Resolving deltas"; + case "checkout": + return "Checking out files"; + } +} + +/** Display name for a clone: the looked-up `owner/repo`, else the folder being cloned into. */ +export function projectCloneDisplayName( + snapshot: Pick, +): string { + if (snapshot.repository) return snapshot.repository.nameWithOwner; + const segments = snapshot.destinationPath.split(/[/\\]/).filter((segment) => segment.length > 0); + return segments[segments.length - 1] ?? snapshot.destinationPath; +} + +/** One-line progress summary: `Receiving objects · 45% · 12.3 MiB | 5.0 MiB/s`. */ +export function projectCloneProgressSummary( + snapshot: Pick, +): string { + const parts = [projectCloneStageLabel(snapshot.stage)]; + if (snapshot.percent !== null) parts.push(`${snapshot.percent}%`); + if (snapshot.detail) parts.push(snapshot.detail); + return parts.join(" · "); +} diff --git a/packages/contracts/src/rpc.ts b/packages/contracts/src/rpc.ts index d792a31885c1..9096dd6e27f9 100644 --- a/packages/contracts/src/rpc.ts +++ b/packages/contracts/src/rpc.ts @@ -249,6 +249,14 @@ import { } from "./providerUsageLimits.ts"; import { UsagePricing, UsageReadError, UsageSummary, UsageSummaryInput } from "./usage.ts"; import { ServerSettings, ServerSettingsError, ServerSettingsPatch } from "./settings.ts"; +import { + ProjectCloneActionInput, + ProjectCloneActionResult, + ProjectCloneListEvent, + ProjectCloneStartInput, + ProjectCloneStartResult, + ProjectCloneSubscribeInput, +} from "./projectClone.ts"; import { SourceControlCloneRepositoryInput, SourceControlCloneRepositoryResult, @@ -407,6 +415,10 @@ export const WS_METHODS = { sourceControlLookupRepository: "sourceControl.lookupRepository", sourceControlCloneRepository: "sourceControl.cloneRepository", sourceControlPublishRepository: "sourceControl.publishRepository", + projectCloneStart: "projectClone.start", + projectCloneCancel: "projectClone.cancel", + projectCloneRetry: "projectClone.retry", + subscribeProjectClones: "subscribeProjectClones", // Streaming subscriptions subscribeVcsStatus: "subscribeVcsStatus", @@ -842,6 +854,37 @@ const WsSourceControlCloneRepositoryRpc = Rpc.make(WS_METHODS.sourceControlClone error: Schema.Union([SourceControlRepositoryError, EnvironmentAuthorizationError]), }); +// Clone-backed project creation. `start` returns once the project exists and +// the clone is running; progress arrives on the subscription. +const WsProjectCloneStartRpc = Rpc.make(WS_METHODS.projectCloneStart, { + payload: ProjectCloneStartInput, + success: ProjectCloneStartResult, + error: Schema.Union([ + SourceControlRepositoryError, + OrchestrationDispatchCommandError, + EnvironmentAuthorizationError, + ]), +}); + +const WsProjectCloneCancelRpc = Rpc.make(WS_METHODS.projectCloneCancel, { + payload: ProjectCloneActionInput, + success: ProjectCloneActionResult, + error: EnvironmentAuthorizationError, +}); + +const WsProjectCloneRetryRpc = Rpc.make(WS_METHODS.projectCloneRetry, { + payload: ProjectCloneActionInput, + success: ProjectCloneActionResult, + error: Schema.Union([SourceControlRepositoryError, EnvironmentAuthorizationError]), +}); + +const WsSubscribeProjectClonesRpc = Rpc.make(WS_METHODS.subscribeProjectClones, { + payload: ProjectCloneSubscribeInput, + success: ProjectCloneListEvent, + error: EnvironmentAuthorizationError, + stream: true, +}); + const WsSourceControlPublishRepositoryRpc = Rpc.make(WS_METHODS.sourceControlPublishRepository, { payload: SourceControlPublishRepositoryInput, success: SourceControlPublishRepositoryResult, @@ -1379,6 +1422,10 @@ export const WsRpcGroup = RpcGroup.make( WsSourceControlLookupRepositoryRpc, WsSourceControlCloneRepositoryRpc, WsSourceControlPublishRepositoryRpc, + WsProjectCloneStartRpc, + WsProjectCloneCancelRpc, + WsProjectCloneRetryRpc, + WsSubscribeProjectClonesRpc, WsProjectsListEntriesRpc, WsProjectsReadFileRpc, WsProjectsSearchContentsRpc,