diff --git a/.plans/21-roaming-workspace.md b/.plans/21-roaming-workspace.md index c6ed9b6fb213..ea527852b21e 100644 --- a/.plans/21-roaming-workspace.md +++ b/.plans/21-roaming-workspace.md @@ -1,6 +1,7 @@ # Roaming workspace — local-first plan -> **Status:** M0 complete 2026-07-04 — both spikes GO (see [M0 results](#m0-results-2026-07-04)); M1 next. +> **Status:** M1 complete 2026-07-04 — exit criteria pass on the harness +> (`scripts/roaming/accept-m1.mjs`; see [M1 results](#m1-results-2026-07-04)); M2 next. > **Decisions log:** 2026-07-04 — v1 transport for small state = machine-to-machine > mirror (user decision); cloud store backend (private git repo or T3 relay) > deferred to explicit milestone M7 behind the same interface. @@ -575,6 +576,44 @@ the plan says. connectivity suffices: mirror RPCs reconcile manifests both ways per contact, so A→B credentials give bidirectional data flow. +## M1 results (2026-07-04) + +Landed as four reviewed PRs into `feature/roaming`: contracts (#2), blob +store (#3), enrollment + PeerMirror (#4), client shell/UI (#5). Exit +criteria verified end-to-end by `scripts/roaming/accept-m1.mjs` on the M0 +harness: enroll on A → registry entry (title, repository, per-machine root) +on B after a mirror pass → A killed → B still serves its local copy. + +Decisions/deviations recorded during implementation and review: + +- **Enrollment is administrative.** The enrollment/mint HTTP routes require + `access:write`, not `orchestration:operate` (any standard client could + otherwise mint 365-day mirror credentials). Consequently the D4 handshake + needs an admin-scoped pairing credential: `t3 auth pairing create --admin` + (new flag). Mirror RPCs require `roaming:mirror`, which is granted nowhere + by default. +- **Peer records are tamper-resistant.** A caller of the machine-credential + route is recorded insert-only (`RoamingPeers.ensurePeer`) and its + advertised base URLs are ignored — overwriting a credentialed peer's URLs + would have redirected our authenticated mirror traffic to an attacker. +- **Enroll ordering:** `project.roaming.enroll` dispatches before the + registry blob write, so the decider gates concurrent double-enrolls and no + orphan blob can mirror out; the idempotent re-enroll path self-heals a + missing blob. +- **Honest staleness:** `last_contact_at` is written only by a completed + mirror pass. Accepted M1 simplification: `lastMirrorContactAt` in the + shell is a global max across peers, not per-project (revisit ~M4). +- **Flag-off behavior:** shell snapshots hide `roamingProjects` and the live + roaming stream while the `roaming` setting is off, consistent with the + routes 404ing. Local blob data is retained. +- **Roaming shell stream events carry `sequence: 0`** (they ride outside the + event log); the client reducer applies them by key and owns all sequencing + rules (the redundant outer gate in shell sync was removed). +- **Reconciliation hardening from review:** the blob store serializes its + read-modify-write behind a semaphore (fiber interleaving could defeat + equal-version conflict detection) and verifies ingested `contentHash` + against the payload before applying. + ## Execution process **Branching (fork discipline):** `main` tracks upstream and receives their diff --git a/apps/web/src/components/Sidebar.tsx b/apps/web/src/components/Sidebar.tsx index 21525b56b770..403b05a87f68 100644 --- a/apps/web/src/components/Sidebar.tsx +++ b/apps/web/src/components/Sidebar.tsx @@ -82,6 +82,7 @@ import { readThreadShell, useProject, useProjects, + useRoamingProjects, useServerConfigs, useThreadShells, useThreadShellsForProjectRefs, @@ -2831,6 +2832,61 @@ const SidebarChromeFooter = memo(function SidebarChromeFooter() { ); }); +/** + * Registry entries mirrored from other machines that are not materialized + * here. Read-only in M1 (materialize arrives with M2); rendered greyed with + * an honest staleness label — the data is only as fresh as the last mirror + * contact. + */ +function SidebarRoamingProjects() { + const roamingProjects = useRoamingProjects(); + const remoteOnly = roamingProjects.filter( + (entry) => entry.roamingProject.localProjectId === null, + ); + if (remoteOnly.length === 0) { + return null; + } + return ( + +
+ + Roaming + +
+ + {remoteOnly.map(({ environmentId, roamingProject }) => { + const repository = + roamingProject.repository.displayName ?? + roamingProject.repository.name ?? + roamingProject.repository.locator.remoteUrl; + const staleness = + roamingProject.lastMirrorContactAt === null + ? "never synced" + : `synced ${formatRelativeTimeLabel(roamingProject.lastMirrorContactAt)}`; + return ( + +
+ +
+
+ {roamingProject.title} +
+
+ {repository} · {staleness} +
+
+
+
+ ); + })} +
+
+ ); +} + interface SidebarProjectsContentProps { showArm64IntelBuildWarning: boolean; arm64IntelBuildWarningDescription: string | null; @@ -3098,6 +3154,7 @@ const SidebarProjectsContent = memo(function SidebarProjectsContent( )} + ); }); diff --git a/apps/web/src/state/entities.ts b/apps/web/src/state/entities.ts index b4fc8cc5e80a..96a87f56add4 100644 --- a/apps/web/src/state/entities.ts +++ b/apps/web/src/state/entities.ts @@ -19,6 +19,7 @@ import { Atom } from "effect/unstable/reactivity"; import { useMemo } from "react"; import { appAtomRegistry } from "../rpc/atomRegistry"; import { environmentProjects } from "./projects"; +import type { EnvironmentRoamingProject } from "@t3tools/client-runtime/state/projects"; import { environmentServerConfigsAtom } from "./server"; import { environmentThreadDetails, environmentThreadShells } from "./threads"; @@ -105,6 +106,10 @@ export function useProjects(): ReadonlyArray { return useAtomValue(environmentProjects.projectsAtom); } +export function useRoamingProjects(): ReadonlyArray { + return useAtomValue(environmentProjects.roamingProjectsAtom); +} + export function useServerConfigs(): ReadonlyMap { return useAtomValue(environmentServerConfigsAtom); } diff --git a/packages/client-runtime/src/state/projectEntities.ts b/packages/client-runtime/src/state/projectEntities.ts index 4d51b4d427e9..b5fa0ac7b8c5 100644 --- a/packages/client-runtime/src/state/projectEntities.ts +++ b/packages/client-runtime/src/state/projectEntities.ts @@ -3,6 +3,7 @@ import type { OrchestrationProjectShell, OrchestrationShellSnapshot, ProjectId, + RoamingProjectShell, ScopedProjectRef, } from "@t3tools/contracts"; import { Atom } from "effect/unstable/reactivity"; @@ -13,6 +14,7 @@ import type { EnvironmentCatalogState } from "./connections.ts"; import { arrayElementsEqual, parseProjectKey, projectKey, projectRefsEqual } from "./entities.ts"; const EMPTY_PROJECTS: ReadonlyArray = Object.freeze([]); +const EMPTY_ROAMING_PROJECTS: ReadonlyArray = Object.freeze([]); const EMPTY_PROJECT_INDEX: ReadonlyMap = new Map(); export function createEnvironmentProjectAtoms(input: { @@ -94,6 +96,35 @@ export function createEnvironmentProjectAtoms(input: { return previousProjects; }).pipe(Atom.withLabel("environment-project-list")); + const environmentRoamingProjectsAtom = Atom.family((environmentId: EnvironmentId) => + Atom.make( + (get): ReadonlyArray => + get(input.snapshotAtom(environmentId))?.roamingProjects ?? EMPTY_ROAMING_PROJECTS, + ).pipe(Atom.withLabel(`environment-roaming-projects:${environmentId}`)), + ); + + let previousRoamingProjects: ReadonlyArray = []; + const roamingProjectsAtom = Atom.make((get) => { + const next: EnvironmentRoamingProject[] = []; + for (const environmentId of get(input.catalogValueAtom).entries.keys()) { + for (const roamingProject of get(environmentRoamingProjectsAtom(environmentId))) { + next.push({ environmentId, roamingProject }); + } + } + const unchanged = + previousRoamingProjects.length === next.length && + next.every( + (entry, index) => + previousRoamingProjects[index]?.environmentId === entry.environmentId && + previousRoamingProjects[index]?.roamingProject === entry.roamingProject, + ); + if (unchanged) { + return previousRoamingProjects; + } + previousRoamingProjects = next; + return previousRoamingProjects; + }).pipe(Atom.withLabel("environment-roaming-project-list")); + return { environmentProjectsAtom, environmentProjectIndexAtom, @@ -101,5 +132,12 @@ export function createEnvironmentProjectAtoms(input: { projectRefsAtom, projectsAtom, projectAtom: (ref: ScopedProjectRef) => projectAtomFamily(projectKey(ref)), + environmentRoamingProjectsAtom, + roamingProjectsAtom, }; } + +export interface EnvironmentRoamingProject { + readonly environmentId: EnvironmentId; + readonly roamingProject: RoamingProjectShell; +} diff --git a/packages/client-runtime/src/state/shell.ts b/packages/client-runtime/src/state/shell.ts index 2b0ba6346f5d..90ce7b8f2117 100644 --- a/packages/client-runtime/src/state/shell.ts +++ b/packages/client-runtime/src/state/shell.ts @@ -130,10 +130,10 @@ export const makeEnvironmentShellState = Effect.fn("EnvironmentShellState.make") ? item.snapshot : Option.match(current.snapshot, { onNone: () => null, - onSome: (snapshot) => - item.sequence > snapshot.snapshotSequence - ? applyShellStreamEvent(snapshot, item) - : snapshot, + // Sequencing rules live in the reducer (roaming events ride + // sequence 0 and are applied by key; sequenced events are + // ignored when stale and return the same reference). + onSome: (snapshot) => applyShellStreamEvent(snapshot, item), }); if (nextSnapshot === null) { return; diff --git a/packages/client-runtime/src/state/shellReducer.test.ts b/packages/client-runtime/src/state/shellReducer.test.ts index 643c25b7fbaa..d6c25bebdb7d 100644 --- a/packages/client-runtime/src/state/shellReducer.test.ts +++ b/packages/client-runtime/src/state/shellReducer.test.ts @@ -1,6 +1,6 @@ import { describe, expect, it } from "vite-plus/test"; -import { ProjectId, ProviderInstanceId, ThreadId } from "@t3tools/contracts"; +import { EnvironmentId, ProjectId, ProviderInstanceId, ThreadId, WorkspaceProjectId } from "@t3tools/contracts"; import type { OrchestrationShellSnapshot, OrchestrationShellStreamEvent } from "@t3tools/contracts"; import { applyShellStreamEvent } from "./shellReducer.ts"; @@ -44,7 +44,58 @@ const stubThread = { session: null, } as const; +const stubRoamingProject = { + workspaceProjectId: WorkspaceProjectId.make("wp-1"), + title: "Roaming Project", + repository: { + canonicalKey: "github.com/acme/app", + locator: { + source: "git-remote" as const, + remoteName: "origin", + remoteUrl: "git@github.com:acme/app.git", + }, + }, + localProjectId: null, + authorEnvironmentId: EnvironmentId.make("env-desktop"), + perMachineRoots: {}, + lastMirrorContactAt: null, + updatedAt: "2026-04-01T00:00:00.000Z", +} as const; + describe("applyShellStreamEvent", () => { + it("applies roaming upserts and removals by key without touching snapshotSequence", () => { + const withHighSequence: OrchestrationShellSnapshot = { + ...baseSnapshot, + snapshotSequence: 10, + }; + + // Roaming events carry sequence 0 and must not be dropped by the guard. + const upserted = applyShellStreamEvent(withHighSequence, { + kind: "roaming-project-upserted", + sequence: 0, + roamingProject: stubRoamingProject, + }); + expect(upserted.roamingProjects).toEqual([stubRoamingProject]); + expect(upserted.snapshotSequence).toBe(10); + + const replaced = applyShellStreamEvent(upserted, { + kind: "roaming-project-upserted", + sequence: 0, + roamingProject: { ...stubRoamingProject, title: "Renamed" }, + }); + expect(replaced.roamingProjects).toHaveLength(1); + expect(replaced.roamingProjects[0]?.title).toBe("Renamed"); + + const removed = applyShellStreamEvent(replaced, { + kind: "roaming-project-removed", + sequence: 0, + workspaceProjectId: stubRoamingProject.workspaceProjectId, + }); + expect(removed.roamingProjects).toEqual([]); + expect(removed.snapshotSequence).toBe(10); + }); + + it("ignores stale project upserts without mutating the snapshot", () => { const snapshotWithProject: OrchestrationShellSnapshot = { ...baseSnapshot, diff --git a/packages/client-runtime/src/state/shellReducer.ts b/packages/client-runtime/src/state/shellReducer.ts index 3d3b22a1289f..fc4ddf42a258 100644 --- a/packages/client-runtime/src/state/shellReducer.ts +++ b/packages/client-runtime/src/state/shellReducer.ts @@ -13,6 +13,33 @@ export function applyShellStreamEvent( snapshot: OrchestrationShellSnapshot, event: OrchestrationShellStreamEvent, ): OrchestrationShellSnapshot { + // Roaming events ride outside the event-log sequence (the server emits + // them with sequence 0): apply by key and leave snapshotSequence alone. + if (event.kind === "roaming-project-upserted") { + const exists = snapshot.roamingProjects.some( + (entry) => entry.workspaceProjectId === event.roamingProject.workspaceProjectId, + ); + return { + ...snapshot, + roamingProjects: exists + ? Arr.map(snapshot.roamingProjects, (entry) => + entry.workspaceProjectId === event.roamingProject.workspaceProjectId + ? event.roamingProject + : entry, + ) + : Arr.append(snapshot.roamingProjects, event.roamingProject), + }; + } + if (event.kind === "roaming-project-removed") { + return { + ...snapshot, + roamingProjects: Arr.filter( + snapshot.roamingProjects, + (entry) => entry.workspaceProjectId !== event.workspaceProjectId, + ), + }; + } + if (event.sequence <= snapshot.snapshotSequence) return snapshot; switch (event.kind) { diff --git a/scripts/roaming/accept-m1.mjs b/scripts/roaming/accept-m1.mjs new file mode 100644 index 000000000000..6407830f2b19 --- /dev/null +++ b/scripts/roaming/accept-m1.mjs @@ -0,0 +1,192 @@ +// M1 acceptance (.plans/21-roaming-workspace.md): enroll a project on +// instance A; instance B holds it (title, repo, per-machine roots) after a +// mirror pass; kill A; B still serves it from its local copy. +// +// Run the harness first: +// scripts/roaming/harness.sh start +// node scripts/roaming/accept-m1.mjs +// +// The script drives only production surfaces: the auth CLI, the orchestration +// dispatch endpoint, and the roaming enrollment + mirror HTTP RPCs. It reads +// A's stored machine credential from A's secret store (local test box) to +// query B the way A's mirror does. + +import { execFileSync } from "node:child_process"; +import { mkdirSync, readFileSync, rmSync, writeFileSync, existsSync } from "node:fs"; +import { randomUUID } from "node:crypto"; +import { join } from "node:path"; + +const HARNESS_DIR = process.env.T3_ROAMING_HARNESS_DIR ?? "/tmp/t3-roaming-harness"; +const A = { name: "instance-a", url: "http://127.0.0.1:14801", base: join(HARNESS_DIR, "instance-a/basedir") }; +const B = { name: "instance-b", url: "http://127.0.0.1:14802", base: join(HARNESS_DIR, "instance-b/basedir") }; +const REPO_ROOT = new URL("../..", import.meta.url).pathname; + +const fail = (step, detail) => { + console.error(`FAIL at ${step}: ${detail}`); + process.exit(1); +}; +const pass = (step) => console.log(`PASS ${step}`); + +const cli = (args) => + execFileSync("node", [join(REPO_ROOT, "apps/server/src/bin.ts"), ...args], { + encoding: "utf8", + }).trim(); + +const api = async (base, path, { method = "GET", token, body } = {}) => { + const response = await fetch(`${base}${path}`, { + method, + headers: { + ...(token ? { authorization: `Bearer ${token}` } : {}), + ...(body ? { "content-type": "application/json" } : {}), + }, + ...(body ? { body: JSON.stringify(body) } : {}), + }); + return response; +}; + +// ── 0. preconditions ───────────────────────────────────────────────── +for (const inst of [A, B]) { + const up = await api(inst.url, "/.well-known/t3/environment").then( + (r) => r.ok, + () => false, + ); + if (!up) fail("preflight", `${inst.name} is not running — start the harness first`); +} + +// ── 1. test repository with an origin remote ───────────────────────── +const workDir = join(HARNESS_DIR, "m1-project"); +const originDir = join(HARNESS_DIR, "m1-origin.git"); +rmSync(workDir, { recursive: true, force: true }); +rmSync(originDir, { recursive: true, force: true }); +execFileSync("git", ["init", "--bare", originDir]); +execFileSync("git", ["init", workDir]); +writeFileSync(join(workDir, "README.md"), "m1 acceptance\n"); +const gitEnv = { + ...process.env, + GIT_AUTHOR_NAME: "m1", + GIT_AUTHOR_EMAIL: "m1@test", + GIT_COMMITTER_NAME: "m1", + GIT_COMMITTER_EMAIL: "m1@test", +}; +execFileSync("git", ["-C", workDir, "add", "."], { env: gitEnv }); +execFileSync("git", ["-C", workDir, "commit", "-m", "init"], { env: gitEnv }); +execFileSync("git", ["-C", workDir, "remote", "add", "origin", originDir]); +pass("test repo prepared"); + +// ── 2. admin bearer on A, project create + enroll ───────────────────── +const adminA = cli(["auth", "session", "issue", "--base-dir", A.base, "--token-only"]); + +const projectId = `m1-${randomUUID()}`; +const dispatchResponse = await api(A.url, "/api/orchestration/dispatch", { + method: "POST", + token: adminA, + body: { + type: "project.create", + commandId: randomUUID(), + projectId, + title: "M1 Acceptance Project", + workspaceRoot: workDir, + createdAt: new Date().toISOString(), + }, +}); +if (!dispatchResponse.ok) fail("project.create", `${dispatchResponse.status} ${await dispatchResponse.text()}`); +pass("project created on A"); + +const enrollResponse = await api(A.url, "/api/roaming/projects/enroll", { + method: "POST", + token: adminA, + body: { projectId }, +}); +if (!enrollResponse.ok) fail("enroll", `${enrollResponse.status} ${await enrollResponse.text()}`); +const { workspaceProjectId } = await enrollResponse.json(); +if (!workspaceProjectId) fail("enroll", "no workspaceProjectId in response"); +pass(`project enrolled on A as ${workspaceProjectId}`); + +// ── 3. peer enrollment A → B (admin pairing credential from B) ──────── +const pairingJson = JSON.parse( + cli(["auth", "pairing", "create", "--admin", "--base-dir", B.base, "--json"]), +); +const pairingCredential = pairingJson.credential ?? pairingJson.token; +if (!pairingCredential) fail("pairing", `unrecognized pairing output: ${JSON.stringify(pairingJson)}`); + +const addPeerResponse = await api(A.url, "/api/roaming/peers", { + method: "POST", + token: adminA, + body: { baseUrls: [B.url], pairingCredential }, +}); +if (!addPeerResponse.ok) fail("addPeer", `${addPeerResponse.status} ${await addPeerResponse.text()}`); +const { peer } = await addPeerResponse.json(); +pass(`peer enrolled: A now mirrors to ${peer.environmentId}`); + +// ── 4. B holds the registry entry after a mirror pass ───────────────── +const credentialPath = join(A.base, "userdata", "secrets", `roaming-peer-${peer.environmentId}.bin`); +if (!existsSync(credentialPath)) fail("credential", `machine credential not stored at ${credentialPath}`); +const machineToken = readFileSync(credentialPath, "utf8"); + +const environmentIdA = readFileSync(join(A.base, "userdata", "environment-id"), "utf8").trim(); + +const waitForBlobOnB = async (timeoutMs) => { + const deadline = Date.now() + timeoutMs; + while (Date.now() < deadline) { + const manifestResponse = await api(B.url, "/api/roaming/mirror/manifest", { + method: "POST", + token: machineToken, + body: { environmentId: environmentIdA, manifest: [] }, + }); + if (manifestResponse.ok) { + const { manifest } = await manifestResponse.json(); + if (manifest.some((entry) => entry.kind === "registry" && entry.key === workspaceProjectId)) { + return true; + } + } + await new Promise((resolve) => setTimeout(resolve, 1000)); + } + return false; +}; + +if (!(await waitForBlobOnB(30_000))) fail("mirror", "registry blob did not reach B within 30s"); +pass("registry entry reached B via mirror pass"); + +const fetchFromB = async () => { + const response = await api(B.url, "/api/roaming/mirror/fetch", { + method: "POST", + token: machineToken, + body: { refs: [{ kind: "registry", key: workspaceProjectId }] }, + }); + if (!response.ok) return null; + const { blobs } = await response.json(); + return blobs[0] ?? null; +}; + +const blob = await fetchFromB(); +if (!blob) fail("fetch", "B did not return the registry blob"); +const payload = JSON.parse(blob.payload); +if (payload.title !== "M1 Acceptance Project") fail("payload", `title mismatch: ${payload.title}`); +if (!payload.repository?.locator?.remoteUrl?.includes("m1-origin.git")) + fail("payload", `repository mismatch: ${JSON.stringify(payload.repository)}`); +if (payload.perMachineRoots?.[environmentIdA] !== workDir) + fail("payload", `perMachineRoots mismatch: ${JSON.stringify(payload.perMachineRoots)}`); +pass("B's copy carries title, repository, and per-machine root"); + +// ── 5. kill A; B still serves its local copy ────────────────────────── +const pidA = Number(readFileSync(join(HARNESS_DIR, "instance-a/server.pid"), "utf8").trim()); +process.kill(pidA, "SIGKILL"); +await new Promise((resolve) => setTimeout(resolve, 1000)); + +const aDown = await api(A.url, "/.well-known/t3/environment").then( + (r) => !r.ok, + () => true, +); +if (!aDown) fail("kill-a", "A is still up"); +const bStillUp = await api(B.url, "/.well-known/t3/environment").then( + (r) => r.ok, + () => false, +); +if (!bStillUp) fail("kill-a", "B went down with A"); + +const blobAfterAKilled = await fetchFromB(); +if (!blobAfterAKilled) fail("survival", "B lost the registry blob after A was killed"); +if (blobAfterAKilled.contentHash !== blob.contentHash) fail("survival", "content hash changed"); +pass("A killed; B still serves the registry entry from its local copy"); + +console.log("\nM1 ACCEPTANCE: ALL CRITERIA PASS");