diff --git a/docs/concepts/workspaces.mdx b/docs/concepts/workspaces.mdx index bbf56428b..5befe8b02 100644 --- a/docs/concepts/workspaces.mdx +++ b/docs/concepts/workspaces.mdx @@ -32,4 +32,4 @@ Think of it as the place where the work happens: you keep the actual source mate ## Source of Truth -The app stores workspace access, the item tree, relationships, current document checkpoints, and extracted text in Postgres. Original and preview file bytes live in R2. The `WorkspaceKernel` Durable Object is a live room for presence, revision notifications, and cleanup coordination; it is not a second database. Clients refetch the authoritative workspace query when a newer revision is announced. UI and AI operations go through workspace commands rather than writing directly to scattered client state. +The app stores workspace access, the item tree, relationships, current document checkpoints, and extracted text in Postgres. Original and preview file bytes live in R2. The `WorkspaceKernel` Durable Object is a live room for presence, workspace-page deltas, and cleanup coordination; it is not a second database. Clients apply canonical item deltas and refetch the authoritative workspace query after reconnecting or when a full refresh is announced. UI and AI operations go through workspace commands rather than writing directly to scattered client state. diff --git a/package.json b/package.json index 568797a61..d37db5c40 100644 --- a/package.json +++ b/package.json @@ -78,7 +78,6 @@ "@posthog/rollup-plugin": "^1.4.7", "@streamdown/cjk": "^1.0.3", "@streamdown/math": "^1.0.2", - "@tanstack/pacer": "^0.22.0", "@tanstack/react-hotkeys": "^0.10.0", "@tanstack/react-query": "5.101.4", "@tanstack/react-router": "1.170.23", diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index 22882b100..b6dfb44f3 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -134,9 +134,6 @@ importers: '@streamdown/math': specifier: ^1.0.2 version: 1.0.2(react@19.2.8) - '@tanstack/pacer': - specifier: ^0.22.0 - version: 0.22.0 '@tanstack/react-hotkeys': specifier: ^0.10.0 version: 0.10.0(react-dom@19.2.8(react@19.2.8))(react@19.2.8) @@ -3638,10 +3635,6 @@ packages: resolution: {integrity: sha512-vqH7X9nb0MTJ/O08++dB5bP9jgj4+BIPOUu/U+6myG86lDsirZSVSobpq5UQpE7nBuk62i8eIYeOhd+OMl/UrA==} engines: {node: '>=18'} - '@tanstack/pacer@0.22.0': - resolution: {integrity: sha512-Zq23KOn30tFiA6ZYWihs5x++ttHkOIcnIIjRrJq05tGrYx9+iHsmVaj8jb7UJVemEI+JRWGzJ9C+50YV4nRz0w==} - engines: {node: '>=18'} - '@tanstack/query-core@5.101.4': resolution: {integrity: sha512-gNwcvOJcRbLWPOLG/2OBm+zM+Yv+MKsXKEOWC57USuZDEsI71hEErQsiEGx5wX9rzWWkfwM0fVSPoiIFSsxfiw==} @@ -11903,11 +11896,6 @@ snapshots: dependencies: '@tanstack/store': 0.11.1 - '@tanstack/pacer@0.22.0': - dependencies: - '@tanstack/devtools-event-client': 0.5.0 - '@tanstack/store': 0.11.1 - '@tanstack/query-core@5.101.4': {} '@tanstack/query-devtools@5.101.4': {} diff --git a/src/features/workspaces/ai/ai-thread-runtime.ts b/src/features/workspaces/ai/ai-thread-runtime.ts index b0134c73c..9897f94a7 100644 --- a/src/features/workspaces/ai/ai-thread-runtime.ts +++ b/src/features/workspaces/ai/ai-thread-runtime.ts @@ -325,6 +325,17 @@ function getWorkspaceAiLanguageModelForGatewayModel( function getWorkspaceAiGatewayTransportOptions() { return { caching: "auto" as const, + // Buy the fast lane where it exists. The gateway only forwards a tier to + // OpenAI, Google AI Studio, and Vertex, so this is a no-op on the Claude + // primaries and moves the models we actually default to (`auto`/luna, the + // Gemini pair, the nano/flash-lite title legs). It is a hint, never a + // promise: an unsupported model ignores it, and a provider that is out of + // priority capacity silently downgrades to standard and bills standard. + // Priority runs ~1.8-2x standard token price when it *is* granted, which + // the `cost` ladder in models.ts does not account for — read + // `service_tier` in PostHog before trusting those multipliers, since a + // missing value means we paid standard and got standard. + serviceTier: "priority" as const, // Time-to-first-token budget before a BYOK leg is abandoned for the next // provider — and eventually for Vercel's own credits. A flat 8s evicted // healthy requests: 37% of gemini-3.1-pro steps and 29% of sonnet steps diff --git a/src/features/workspaces/cache-page.test.ts b/src/features/workspaces/cache-page.test.ts new file mode 100644 index 000000000..fc9b24ea5 --- /dev/null +++ b/src/features/workspaces/cache-page.test.ts @@ -0,0 +1,87 @@ +import { QueryClient } from "@tanstack/react-query"; +import { describe, expect, it, vi } from "vitest"; + +import { workspacePageQueryKey } from "#/features/workspaces/cache-keys"; +import { applyWorkspacePageDeltaToCache } from "#/features/workspaces/cache-page"; +import type { WorkspaceItemSummary, WorkspacePage } from "#/features/workspaces/contracts"; + +describe("workspace page cache ordering", () => { + it("applies only the next revision", () => { + const queryClient = createQueryClient(createPage(3, createItem({ name: "Before" }))); + + applyWorkspacePageDeltaToCache(queryClient, { + type: "workspace.items.upserted", + workspaceId: "workspace-1", + revision: 4, + items: [createItem({ name: "After" })], + }); + + expect(readPage(queryClient)).toMatchObject({ + items: [{ name: "After" }], + revision: 4, + }); + }); + + it("ignores stale revisions", () => { + const queryClient = createQueryClient(createPage(3, createItem({ name: "Current" }))); + const invalidate = vi.spyOn(queryClient, "invalidateQueries").mockResolvedValue(); + + applyWorkspacePageDeltaToCache(queryClient, { + type: "workspace.items.upserted", + workspaceId: "workspace-1", + revision: 2, + items: [createItem({ name: "Stale" })], + }); + + expect(readItem(queryClient)).toMatchObject({ name: "Current" }); + expect(invalidate).not.toHaveBeenCalled(); + }); + + it("keeps current cache data and reconciles a revision gap", () => { + const queryClient = createQueryClient(createPage(3, createItem({ name: "Newer" }))); + const invalidate = vi.spyOn(queryClient, "invalidateQueries").mockResolvedValue(); + + applyWorkspacePageDeltaToCache(queryClient, { + type: "workspace.items.upserted", + workspaceId: "workspace-1", + revision: 5, + items: [createItem({ name: "Older" })], + }); + + expect(readItem(queryClient)).toMatchObject({ name: "Newer" }); + expect(invalidate).toHaveBeenCalledWith({ queryKey: workspacePageQueryKey("workspace-1") }); + }); +}); + +function createQueryClient(page: WorkspacePage) { + const queryClient = new QueryClient(); + queryClient.setQueryData(workspacePageQueryKey("workspace-1"), page); + return queryClient; +} + +function readItem(queryClient: QueryClient) { + return readPage(queryClient)?.items[0]; +} + +function readPage(queryClient: QueryClient) { + return queryClient.getQueryData(workspacePageQueryKey("workspace-1")); +} + +function createPage(revision: number, item: WorkspaceItemSummary): WorkspacePage { + return { workspace: {} as WorkspacePage["workspace"], items: [item], revision }; +} + +function createItem(input: Partial = {}): WorkspaceItemSummary { + return { + color: input.color ?? null, + createdAt: "2026-01-01T00:00:00.000Z", + id: "folder-1", + metadataJson: {}, + name: input.name ?? "Folder", + parentId: null, + sortOrder: 1, + type: "folder", + updatedAt: "2026-01-01T00:00:00.000Z", + workspaceId: "workspace-1", + }; +} diff --git a/src/features/workspaces/cache-page.ts b/src/features/workspaces/cache-page.ts index 744f48201..dfcf4165f 100644 --- a/src/features/workspaces/cache-page.ts +++ b/src/features/workspaces/cache-page.ts @@ -3,15 +3,40 @@ import { workspacePageQueryKey } from "#/features/workspaces/cache-keys"; import type { CreateWorkspaceItemInput, MoveWorkspaceItemsInput, - UpdateWorkspaceItemColorInput, WorkspacePage, } from "#/features/workspaces/contracts"; import { createWorkspaceItemInPage, moveWorkspaceItemsInPage, removeWorkspaceItemsFromPage, - updateWorkspaceItemColorInPage, + upsertWorkspaceItemInPage, } from "#/features/workspaces/model/workspace-page"; +import type { WorkspacePageDelta } from "#/features/workspaces/realtime/messages"; + +export function applyWorkspacePageDeltaToCache( + queryClient: QueryClient, + change: WorkspacePageDelta, +) { + let shouldReconcile = false; + queryClient.setQueryData(workspacePageQueryKey(change.workspaceId), (current) => { + if (!current) return current; + if (change.revision <= current.revision) return current; + if (change.revision !== current.revision + 1) { + shouldReconcile = true; + return current; + } + if (change.type === "workspace.items.deleted") { + return removeWorkspaceItemsFromPage(current, change.itemIds, change.revision); + } + return change.items.reduce( + (page, item) => upsertWorkspaceItemInPage(page, item, change.revision), + current, + ); + }); + if (shouldReconcile) { + void queryClient.invalidateQueries({ queryKey: workspacePageQueryKey(change.workspaceId) }); + } +} export function createWorkspaceItemInPageCache( queryClient: QueryClient, @@ -40,31 +65,3 @@ export function removeWorkspaceItemsFromPageCache( current ? removeWorkspaceItemsFromPage(current, itemIds) : current, ); } - -export function updateWorkspaceItemColorInPageCache( - queryClient: QueryClient, - input: UpdateWorkspaceItemColorInput, -) { - queryClient.setQueryData(workspacePageQueryKey(input.workspaceId), (current) => { - if (!current) { - return current; - } - - const updateResult = updateWorkspaceItemColorInPage(current, input); - - if (!updateResult) { - return current; - } - - return updateResult; - }); -} - -export function getWorkspaceItemColorInPageCache( - queryClient: QueryClient, - input: Pick, -) { - const page = queryClient.getQueryData(workspacePageQueryKey(input.workspaceId)); - - return page?.items.find((item) => item.id === input.itemId)?.color ?? null; -} diff --git a/src/features/workspaces/cache-workspace.ts b/src/features/workspaces/cache-workspace.ts index d569f8505..722503ce3 100644 --- a/src/features/workspaces/cache-workspace.ts +++ b/src/features/workspaces/cache-workspace.ts @@ -38,14 +38,12 @@ export function setWorkspacePageCache( input: { workspace: WorkspaceSummary; items: WorkspaceItemSummary[]; - itemFacts: WorkspacePage["itemFacts"]; revision: number; }, ) { queryClient.setQueryData(workspacePageQueryKey(input.workspace.id), { workspace: input.workspace, items: input.items, - itemFacts: input.itemFacts, revision: input.revision, }); } diff --git a/src/features/workspaces/components/WorkspaceFileUploadProvider.tsx b/src/features/workspaces/components/WorkspaceFileUploadProvider.tsx index 7511a8500..3c3d3b2ac 100644 --- a/src/features/workspaces/components/WorkspaceFileUploadProvider.tsx +++ b/src/features/workspaces/components/WorkspaceFileUploadProvider.tsx @@ -14,7 +14,7 @@ import { } from "#/components/ui/alert-dialog"; import { useBillingState } from "#/features/account/use-billing-state"; import { showUpgradeDialog } from "#/features/account/upgrade-navigation"; -import { workspacePageQueryKey } from "#/features/workspaces/cache"; +import { applyWorkspacePageDeltaToCache } from "#/features/workspaces/cache"; import { useWorkspaceMutationAccess } from "#/features/workspaces/components/workspace-mutation-access"; import { runWorkspaceFileUploadBatch } from "#/features/workspaces/files/workspace-file-upload"; import { workspaceUploadAccept } from "#/features/workspaces/upload/workspace-upload-intake"; @@ -57,8 +57,13 @@ export function WorkspaceFileUploadProvider({ parentId, files: fileList, onLimitReached: setLimitResult, - onSuccess: () => { - void queryClient.invalidateQueries({ queryKey: workspacePageQueryKey(workspaceId) }); + onSuccess: (command) => { + applyWorkspacePageDeltaToCache(queryClient, { + type: "workspace.items.upserted", + workspaceId, + items: [command.result], + revision: command.revision, + }); }, }); }; diff --git a/src/features/workspaces/components/WorkspaceItemActionsMenu.tsx b/src/features/workspaces/components/WorkspaceItemActionsMenu.tsx index 6610893f3..93d36bb30 100644 --- a/src/features/workspaces/components/WorkspaceItemActionsMenu.tsx +++ b/src/features/workspaces/components/WorkspaceItemActionsMenu.tsx @@ -113,7 +113,7 @@ export function WorkspaceItemActionsMenuContent({ updateWorkspaceItemColorMutation.mutate({ workspaceId: item.workspaceId, @@ -193,12 +193,12 @@ function WorkspaceItemRenameMenuItem({ function WorkspaceItemColorSubmenu({ item, menuKind, - readOnly, + disabled, onUpdateItemColor, }: { item: WorkspaceItem; menuKind: "dropdown" | "context"; - readOnly: boolean; + disabled: boolean; onUpdateItemColor: (color: WorkspaceItemColor) => void; }) { const selectedColor = getWorkspaceItemColorValue(item.color); @@ -210,14 +210,14 @@ function WorkspaceItemColorSubmenu({ onValueChange={onUpdateItemColor} showLabels={false} className="grid-flow-col grid-rows-4 gap-1.5" - disabled={readOnly} + disabled={disabled} /> ); if (menuKind === "context") { return ( - + {workspaceItemColorSubmenuTrigger} @@ -229,7 +229,7 @@ function WorkspaceItemColorSubmenu({ return ( - + {workspaceItemColorSubmenuTrigger} diff --git a/src/features/workspaces/components/WorkspaceLayout.tsx b/src/features/workspaces/components/WorkspaceLayout.tsx index 71bc894c9..c499bc065 100644 --- a/src/features/workspaces/components/WorkspaceLayout.tsx +++ b/src/features/workspaces/components/WorkspaceLayout.tsx @@ -1,6 +1,6 @@ import { useQueryClient } from "@tanstack/react-query"; import { useEffect } from "react"; -import { workspacePageQueryKey } from "#/features/workspaces/cache"; +import { applyWorkspacePageDeltaToCache, workspacePageQueryKey } from "#/features/workspaces/cache"; import AiChatPanel from "#/features/workspaces/components/AiChatPanel"; import WorkspaceChatLayout from "#/features/workspaces/components/WorkspaceChatLayout"; import WorkspaceContextBar from "#/features/workspaces/components/WorkspaceContextBar"; @@ -22,11 +22,7 @@ import { useWorkspaceViewPolicy, WorkspaceViewCapabilitiesProvider, } from "#/features/workspaces/components/workspace-view-policy"; -import type { - WorkspaceItemFacts, - WorkspaceItemType, - WorkspaceSummary, -} from "#/features/workspaces/contracts"; +import type { WorkspaceItemType, WorkspaceSummary } from "#/features/workspaces/contracts"; import type { WorkspaceLocation } from "#/features/workspaces/locations/workspace-location"; import { WorkspaceLocationProvider } from "#/features/workspaces/locations/workspace-location-context"; import { DocumentEditReviewProvider } from "#/features/workspaces/documents/document-edit-review-context"; @@ -57,8 +53,6 @@ export type { WorkspaceItem } from "#/features/workspaces/model/types"; interface WorkspaceShellProps { workspace: WorkspaceSummary; items: WorkspaceItem[]; - itemFacts: WorkspaceItemFacts[]; - revision: number; activeTabIdFromUrl?: string; activeViewFromUrl?: string; } @@ -66,8 +60,6 @@ interface WorkspaceShellProps { export function WorkspaceShell({ workspace, items, - itemFacts, - revision, activeTabIdFromUrl, activeViewFromUrl, }: WorkspaceShellProps) { @@ -83,8 +75,10 @@ export function WorkspaceShell({ const normalizedUiSession = useWorkspaceUiSession(workspace.id); const realtime = useWorkspaceRealtime({ workspaceId: workspace.id, - lastSeenRevision: revision, - onWorkspaceChanged: () => { + onPageChange: (change) => { + applyWorkspacePageDeltaToCache(queryClient, change); + }, + onDesync: () => { void queryClient.invalidateQueries({ queryKey: workspacePageQueryKey(workspace.id), }); @@ -183,7 +177,6 @@ export function WorkspaceShell({ activeItem: isWorkspaceItemView(activeItem) ? activeItem : undefined, activeTabId: activeTab.id, itemViewStatesByItemId, - itemFactsById: new Map(itemFacts.map((item) => [item.itemId, item])), itemsById, presentation, selectedItemIds, diff --git a/src/features/workspaces/components/WorkspacePageRoute.tsx b/src/features/workspaces/components/WorkspacePageRoute.tsx index 450979b14..d21148cbf 100644 --- a/src/features/workspaces/components/WorkspacePageRoute.tsx +++ b/src/features/workspaces/components/WorkspacePageRoute.tsx @@ -75,8 +75,6 @@ export default function WorkspacePageRoute() { activeTabIdFromUrl={tab} activeViewFromUrl={view} items={page.items} - itemFacts={page.itemFacts} - revision={page.revision} workspace={page.workspace} /> ); diff --git a/src/features/workspaces/contracts.ts b/src/features/workspaces/contracts.ts index f4b3b674c..82f8c03bc 100644 --- a/src/features/workspaces/contracts.ts +++ b/src/features/workspaces/contracts.ts @@ -348,12 +348,6 @@ export const workspaceSummarySchema = z.object({ membershipRole: workspaceMembershipRoleSchema, }); -export const workspaceItemFactsSchema = z.object({ - itemId: z.string(), - pageCount: z.number().int().positive().optional(), - relationshipCount: z.number().int().nonnegative(), -}); - export const workspaceItemSummarySchema = z.object({ id: z.string(), workspaceId: z.string(), @@ -448,7 +442,6 @@ export const workspaceIdInputSchema = z.object({ export const workspacePageSchema = z.object({ workspace: workspaceSummarySchema, items: z.array(workspaceItemSummarySchema), - itemFacts: z.array(workspaceItemFactsSchema), revision: z.number().int().nonnegative(), }); @@ -457,7 +450,6 @@ export type WorkspaceColor = z.infer; export type WorkspaceItemColor = z.infer; export type WorkspaceSummary = z.infer; export type WorkspaceDetail = WorkspaceSummary; -export type WorkspaceItemFacts = z.infer; export type WorkspaceItemSummary = z.infer; export type CreateWorkspaceItemInput = z.infer; export type RenameWorkspaceItemInput = z.infer; diff --git a/src/features/workspaces/kernel/workspace-kernel-access.ts b/src/features/workspaces/kernel/workspace-kernel-access.ts index d54c725a9..2926ad70f 100644 --- a/src/features/workspaces/kernel/workspace-kernel-access.ts +++ b/src/features/workspaces/kernel/workspace-kernel-access.ts @@ -5,7 +5,6 @@ import type { MoveWorkspaceItemsInput, RenameWorkspaceItemInput, UpdateWorkspaceItemColorInput, - WorkspaceItemFacts, WorkspaceItemSummary, } from "#/features/workspaces/contracts"; import { @@ -50,15 +49,12 @@ export interface WorkspaceKernelClient { getPage(input?: { userId?: string }): Promise<{ workspaceId: string; items: WorkspaceItemSummary[]; - itemFacts: WorkspaceItemFacts[]; revision: number; }>; listTreeItems(input?: ListWorkspaceKernelItemsArgs): Promise; resolvePaths(input: ResolveWorkspaceKernelPathsArgs): Promise; getItemPaths(input: GetWorkspaceKernelItemPathsArgs): Promise; - linkItems( - input: LinkWorkspaceKernelItemsArgs, - ): Promise>; + linkItems(input: LinkWorkspaceKernelItemsArgs): Promise; listItemRelations( input: ListWorkspaceKernelItemRelationsArgs, ): Promise; diff --git a/src/features/workspaces/kernel/workspace-kernel-list.ts b/src/features/workspaces/kernel/workspace-kernel-list.ts index 3c4b4f383..0707a9db5 100644 --- a/src/features/workspaces/kernel/workspace-kernel-list.ts +++ b/src/features/workspaces/kernel/workspace-kernel-list.ts @@ -1,8 +1,4 @@ -import type { - WorkspaceItemFacts, - WorkspaceItemSummary, - WorkspaceItemType, -} from "#/features/workspaces/contracts"; +import type { WorkspaceItemSummary, WorkspaceItemType } from "#/features/workspaces/contracts"; import { joinWorkspacePathSegment, resolveWorkspaceKernelCwd, @@ -25,9 +21,7 @@ export interface ListWorkspaceKernelItemsResult { export interface ListWorkspaceKernelItem { modifiedAt: string; - pageCount?: number; path: string; - relationshipCount: number; type: WorkspaceItemType | WorkspaceFileAssetKind; } @@ -50,18 +44,13 @@ interface WorkspaceKernelListRow { } export function listWorkspaceKernelTreeItems(input: { - getItemFacts: (items: WorkspaceItemSummary[]) => WorkspaceItemFacts[]; tree: WorkspaceKernelTree; offset?: number; path?: string; recursive?: boolean; limit?: number; }): ListWorkspaceKernelItemsResult { - const selection = selectWorkspaceKernelTreeItems(input); - return formatWorkspaceKernelListSelection( - selection, - input.getItemFacts(selection.rows.map((row) => row.item)), - ); + return formatWorkspaceKernelListSelection(selectWorkspaceKernelTreeItems(input)); } function selectWorkspaceKernelTreeItems(input: { @@ -112,14 +101,11 @@ function selectWorkspaceKernelTreeItems(input: { function formatWorkspaceKernelListSelection( selection: WorkspaceKernelListSelection, - itemFacts: WorkspaceItemFacts[], ): ListWorkspaceKernelItemsResult { - const itemFactsById = new Map(itemFacts.map((facts) => [facts.itemId, facts])); return { failed: selection.failed, items: selection.rows.map((row) => formatWorkspaceKernelListItem({ - facts: itemFactsById.get(row.item.id), item: row.item, path: row.path, }), @@ -180,15 +166,12 @@ function collectWorkspaceKernelListRows({ } function formatWorkspaceKernelListItem(input: { - facts?: WorkspaceItemFacts; item: WorkspaceItemSummary; path: string; }): ListWorkspaceKernelItem { return { modifiedAt: input.item.updatedAt, - ...(input.facts?.pageCount ? { pageCount: input.facts.pageCount } : {}), path: input.path, - relationshipCount: input.facts?.relationshipCount ?? 0, type: resolveWorkspaceFileTypeFromItem(input.item)?.assetKind ?? input.item.type, }; } diff --git a/src/features/workspaces/kernel/workspace-kernel.ts b/src/features/workspaces/kernel/workspace-kernel.ts index 3c922bb41..644b95d29 100644 --- a/src/features/workspaces/kernel/workspace-kernel.ts +++ b/src/features/workspaces/kernel/workspace-kernel.ts @@ -10,7 +10,7 @@ import { import type { ResourcePurgeResult } from "#/features/workspaces/resource-purge-result"; import type { WorkspaceConnectionState, - WorkspaceRevision, + WorkspacePageChange, WorkspaceRealtimeServerMessage, } from "#/features/workspaces/realtime/messages"; import { recordOperationalFailure } from "#/integrations/observability/operational-events"; @@ -23,7 +23,7 @@ export { setWorkspaceKernelUserHeaders }; /** * The deployed class name is retained to avoid a destructive Durable Object * migration. Canonical workspace data lives in Postgres; this object is only a - * live room for presence, revision notifications, and retryable remote cleanup. + * live room for presence, workspace-page deltas, and retryable remote cleanup. */ export class WorkspaceKernel extends Agent { onConnect(connection: Connection, context: ConnectionContext) { @@ -42,16 +42,12 @@ export class WorkspaceKernel extends Agent { this.broadcastPresenceSnapshot(); } - async publishWorkspaceChange(change: WorkspaceRevision): Promise { + async publishWorkspacePageChange(change: WorkspacePageChange): Promise { if (change.workspaceId !== this.name) { throw new Error("Workspace change was routed to the wrong room."); } - this.broadcastRealtimeMessage({ - type: "workspace.changed", - workspaceId: this.name, - revision: change.revision, - }); + this.broadcastRealtimeMessage(change); } async disconnectMember(input: { userId: string; documentItemIds: string[] }): Promise { diff --git a/src/features/workspaces/model/workspace-ai-context-outline.ts b/src/features/workspaces/model/workspace-ai-context-outline.ts index 9620e7060..f927292f4 100644 --- a/src/features/workspaces/model/workspace-ai-context-outline.ts +++ b/src/features/workspaces/model/workspace-ai-context-outline.ts @@ -59,13 +59,9 @@ function getWorkspaceAiContextOutlineRows( throw new Error("Workspace outline path index returned an unknown item."); } - const facts = context.itemFactsById.get(item.id); - return { item, outlineItem: { - ...(facts?.pageCount ? { pageCount: facts.pageCount } : {}), - relationshipCount: facts?.relationshipCount ?? 0, ...(item.type === "folder" ? { childCount: childCountsByItemId.get(item.id) ?? 0, diff --git a/src/features/workspaces/model/workspace-ai-context-prompt.ts b/src/features/workspaces/model/workspace-ai-context-prompt.ts index fd0f0cfe6..54d459768 100644 --- a/src/features/workspaces/model/workspace-ai-context-prompt.ts +++ b/src/features/workspaces/model/workspace-ai-context-prompt.ts @@ -106,14 +106,7 @@ function formatWorkspaceAiContextOutlineItemMeta(item: WorkspaceAiContextOutline ? "" : `, ${item.childCount} direct ${item.childCount === 1 ? "child" : "children"}, ${item.descendantCount} total ${item.descendantCount === 1 ? "descendant" : "descendants"}`; - const pages = item.pageCount - ? `, ${item.pageCount} ${item.pageCount === 1 ? "page" : "pages"}` - : ""; - const relationships = item.relationshipCount - ? `, ${item.relationshipCount} ${item.relationshipCount === 1 ? "relationship" : "relationships"}` - : ""; - - return `${item.type}${pages}${relationships}${counts}`; + return `${item.type}${counts}`; } function limitWorkspaceAiContextOutlineLines(lines: string[]) { diff --git a/src/features/workspaces/model/workspace-ai-context-types.ts b/src/features/workspaces/model/workspace-ai-context-types.ts index e56fa2bf0..211dfebda 100644 --- a/src/features/workspaces/model/workspace-ai-context-types.ts +++ b/src/features/workspaces/model/workspace-ai-context-types.ts @@ -1,6 +1,5 @@ import type { WorkspaceTab } from "#/features/workspaces/model/tab-types"; import type { WorkspaceItem } from "#/features/workspaces/model/types"; -import type { WorkspaceItemFacts } from "#/features/workspaces/contracts"; import type { WorkspaceAiContextItemViewState, WorkspaceItemViewState, @@ -12,7 +11,6 @@ export type WorkspaceAiContextScope = { activeItem?: WorkspaceItem; activeTabId?: string; itemViewStatesByItemId: Readonly>; - itemFactsById: ReadonlyMap; itemsById: ReadonlyMap; presentation: WorkspacePresentation; selectedItemIds: readonly string[]; @@ -55,9 +53,7 @@ export type WorkspaceAiContextOutline = export type WorkspaceAiContextOutlineItem = { childCount?: number; descendantCount?: number; - pageCount?: number; path: string; - relationshipCount: number; type: string; }; diff --git a/src/features/workspaces/model/workspace-ai-context-validation.ts b/src/features/workspaces/model/workspace-ai-context-validation.ts index be4d4a0ab..dc6141ca1 100644 --- a/src/features/workspaces/model/workspace-ai-context-validation.ts +++ b/src/features/workspaces/model/workspace-ai-context-validation.ts @@ -84,8 +84,6 @@ function isWorkspaceAiContextOutlineItem(value: unknown) { return ( typeof value.path === "string" && typeof value.type === "string" && - (value.pageCount === undefined || isPositiveInteger(value.pageCount)) && - isNonNegativeInteger(value.relationshipCount) && hasChildCount === hasDescendantCount && (value.childCount === undefined || isNonNegativeInteger(value.childCount)) && (value.descendantCount === undefined || isNonNegativeInteger(value.descendantCount)) @@ -96,10 +94,6 @@ function isNonNegativeInteger(value: unknown) { return typeof value === "number" && Number.isInteger(value) && value >= 0; } -function isPositiveInteger(value: unknown) { - return typeof value === "number" && Number.isInteger(value) && value > 0; -} - function isWorkspaceAiContextSelectedItem(value: unknown): value is WorkspaceAiContextSelectedItem { if (!isRecord(value) || !isWorkspaceAiContextItemReference(value)) { return false; diff --git a/src/features/workspaces/model/workspace-page.test.ts b/src/features/workspaces/model/workspace-page.test.ts index f57de50ee..bed656837 100644 --- a/src/features/workspaces/model/workspace-page.test.ts +++ b/src/features/workspaces/model/workspace-page.test.ts @@ -13,18 +13,11 @@ describe("removeWorkspaceItemsFromPage", () => { createItem({ id: "grandchild", parentId: "child" }), createItem({ id: "sibling" }), ], - itemFacts: [ - { itemId: "folder", relationshipCount: 0 }, - { itemId: "child", relationshipCount: 0 }, - { itemId: "grandchild", relationshipCount: 0 }, - { itemId: "sibling", relationshipCount: 0 }, - ], revision: 1, } satisfies WorkspacePage; expect(removeWorkspaceItemsFromPage(page, ["folder"])).toMatchObject({ items: [{ id: "sibling" }], - itemFacts: [{ itemId: "sibling" }], }); }); }); diff --git a/src/features/workspaces/model/workspace-page.ts b/src/features/workspaces/model/workspace-page.ts index 2d236e4b3..74212f118 100644 --- a/src/features/workspaces/model/workspace-page.ts +++ b/src/features/workspaces/model/workspace-page.ts @@ -1,7 +1,6 @@ import type { CreateWorkspaceItemInput, MoveWorkspaceItemsInput, - UpdateWorkspaceItemColorInput, WorkspaceItemSummary, WorkspacePage, } from "#/features/workspaces/contracts"; @@ -116,23 +115,6 @@ export function moveWorkspaceItemsInPage( return nextPage; } -export function updateWorkspaceItemColorInPage( - page: WorkspacePage, - input: UpdateWorkspaceItemColorInput, -): WorkspacePage | null { - const previousItem = page.items.find((item) => item.id === input.itemId); - - if (!previousItem) { - return null; - } - - return upsertWorkspaceItemInPage(page, { - ...previousItem, - color: input.color, - updatedAt: new Date().toISOString(), - }); -} - export function upsertWorkspaceItemInPage( page: WorkspacePage, item: WorkspaceItemSummary, @@ -161,7 +143,6 @@ export function removeWorkspaceItemsFromPage( ...page, revision: Math.max(page.revision, revision), items: page.items.filter((item) => !deletedIds.has(item.id)), - itemFacts: page.itemFacts.filter((facts) => !deletedIds.has(facts.itemId)), }; } diff --git a/src/features/workspaces/operations/workspace-tool-schemas.ts b/src/features/workspaces/operations/workspace-tool-schemas.ts index f1c34a35c..25f88e8e0 100644 --- a/src/features/workspaces/operations/workspace-tool-schemas.ts +++ b/src/features/workspaces/operations/workspace-tool-schemas.ts @@ -71,9 +71,7 @@ const workspacePathItemSchema = z.object({ const workspaceListItemSchema = z.object({ modifiedAt: z.string(), - pageCount: z.number().int().positive().optional(), path: workspacePathSchema, - relationshipCount: z.number().int().nonnegative(), type: z.union([workspaceItemTypeSchema, workspaceFileAssetKindSchema]), }); diff --git a/src/features/workspaces/persistence/workspace-postgres-documents.ts b/src/features/workspaces/persistence/workspace-postgres-documents.ts index b083926e8..f6bc033cd 100644 --- a/src/features/workspaces/persistence/workspace-postgres-documents.ts +++ b/src/features/workspaces/persistence/workspace-postgres-documents.ts @@ -6,7 +6,7 @@ import type { ReadWorkspaceDocumentCheckpointArgs, WorkspaceKernelPublishOutcome, } from "#/features/workspaces/kernel/workspace-kernel-types"; -import type { WorkspaceRevision } from "#/features/workspaces/realtime/messages"; +import type { WorkspacePageDelta } from "#/features/workspaces/realtime/messages"; import { getActiveWorkspaceItemRow, lockWorkspaceForActor, @@ -20,7 +20,7 @@ import { export class PostgresWorkspaceDocuments { constructor( private readonly workspaceId: string, - private readonly onChange?: (change: WorkspaceRevision) => Promise, + private readonly onChange?: (change: WorkspacePageDelta) => Promise, ) {} async readCheckpoint(input: ReadWorkspaceDocumentCheckpointArgs) { @@ -69,11 +69,17 @@ export class PostgresWorkspaceDocuments { .update(workspaceItems) .set({ metadata, updatedAt: now }) .where(eq(workspaceItems.id, input.itemId)); + const item = await requireActiveWorkspaceItem(transaction, this.workspaceId, input.itemId); const revision = await nextWorkspaceRevision(transaction, this.workspaceId); - return { outcome: "applied" as const, revision }; + return { outcome: "applied" as const, item, revision }; }); if (publication.outcome === "applied") { - await this.onChange?.({ workspaceId: this.workspaceId, revision: publication.revision }); + await this.onChange?.({ + type: "workspace.items.upserted", + workspaceId: this.workspaceId, + revision: publication.revision, + items: [publication.item], + }); } return publication.outcome; } diff --git a/src/features/workspaces/persistence/workspace-postgres-files.ts b/src/features/workspaces/persistence/workspace-postgres-files.ts index f9ba59b70..781b9b5a5 100644 --- a/src/features/workspaces/persistence/workspace-postgres-files.ts +++ b/src/features/workspaces/persistence/workspace-postgres-files.ts @@ -24,7 +24,7 @@ import { } from "#/features/workspaces/model/workspace-file"; import type { WorkspaceCommandResult, - WorkspaceRevision, + WorkspacePageDelta, } from "#/features/workspaces/realtime/messages"; import { assertCanReadWorkspace } from "#/features/workspaces/server/permissions"; import { @@ -47,7 +47,7 @@ export class PostgresWorkspaceFiles { constructor( private readonly workspaceId: string, private readonly bucket: R2Bucket, - private readonly onChange?: (change: WorkspaceRevision) => Promise, + private readonly onChange?: (change: WorkspacePageDelta) => Promise, ) {} async createFileFromUpload( @@ -125,7 +125,12 @@ export class PostgresWorkspaceFiles { const revision = await nextWorkspaceRevision(transaction, this.workspaceId); return { result: item, revision }; }); - await this.notify(command.revision); + await this.onChange?.({ + type: "workspace.items.upserted", + workspaceId: this.workspaceId, + revision: command.revision, + items: [command.result], + }); return command; } @@ -223,10 +228,8 @@ export class PostgresWorkspaceFiles { updatedAt: now, }, }); - const revision = await nextWorkspaceRevision(transaction, this.workspaceId); - return { outcome: "applied" as const, revision }; + return { outcome: "applied" as const }; }); - if (publication.outcome === "applied") await this.notify(publication.revision); return publication.outcome; } @@ -310,10 +313,8 @@ export class PostgresWorkspaceFiles { updatedAt: now, }, }); - const revision = await nextWorkspaceRevision(transaction, this.workspaceId); - return { outcome: "applied" as const, revision }; + return { outcome: "applied" as const }; }); - if (publication.outcome === "applied") await this.notify(publication.revision); return publication.outcome; } @@ -337,8 +338,4 @@ export class PostgresWorkspaceFiles { .orderBy(asc(workspaceItemPages.pageNumber)); }); } - - private async notify(revision: number) { - await this.onChange?.({ workspaceId: this.workspaceId, revision }); - } } diff --git a/src/features/workspaces/persistence/workspace-postgres-persistence.ts b/src/features/workspaces/persistence/workspace-postgres-persistence.ts index 7b08abbef..eb22d8508 100644 --- a/src/features/workspaces/persistence/workspace-postgres-persistence.ts +++ b/src/features/workspaces/persistence/workspace-postgres-persistence.ts @@ -42,7 +42,7 @@ import { } from "#/features/workspaces/model/workspace-item-colors"; import type { WorkspaceCommandResult, - WorkspaceRevision, + WorkspacePageDelta, } from "#/features/workspaces/realtime/messages"; import { assertCanReadWorkspace } from "#/features/workspaces/server/permissions"; import { PostgresWorkspaceDocuments } from "./workspace-postgres-documents"; @@ -52,9 +52,9 @@ import { collectDescendants, getActiveWorkspaceItemRows, getNextWorkspaceSortOrder, - getWorkspaceItemFacts, getActiveWorkspaceItemRow, getWorkspaceItemsByIds, + getWorkspaceRevision, hasSelectedAncestor, isDescendantOf, lockWorkspaceForActor, @@ -71,8 +71,8 @@ import { /** * Postgres implementation of the existing workspace-kernel persistence surface. * - * Mutations commit a workspace revision atomically, then publish that revision - * as a cache-invalidation hint through the injected callback. + * Page-affecting mutations commit a workspace revision atomically, then publish + * their canonical cache delta through the injected callback. */ export class PostgresWorkspacePersistence implements WorkspaceKernelClient { private readonly files: PostgresWorkspaceFiles; @@ -81,7 +81,7 @@ export class PostgresWorkspacePersistence implements WorkspaceKernelClient { constructor( private readonly workspaceId: string, bucket: R2Bucket, - private readonly onChange?: (change: WorkspaceRevision) => Promise, + private readonly onChange?: (change: WorkspacePageDelta) => Promise, private readonly onItemsDeleted?: (input: { workspaceId: string; documentItemIds: string[]; @@ -111,12 +111,9 @@ export class PostgresWorkspacePersistence implements WorkspaceKernelClient { async listTreeItems(input: ListWorkspaceKernelItemsArgs = {}) { const page = await this.getPage(); - const factsById = new Map(page.itemFacts.map((facts) => [facts.itemId, facts])); return listWorkspaceKernelTreeItems({ ...input, tree: buildWorkspaceKernelTree(page.items), - getItemFacts: (items) => - items.map((item) => factsById.get(item.id) ?? { itemId: item.id, relationshipCount: 0 }), }); } @@ -178,7 +175,7 @@ export class PostgresWorkspacePersistence implements WorkspaceKernelClient { } async linkItems(input: LinkWorkspaceKernelItemsArgs) { - const command = await withWorkspaceTransaction(async (transaction) => { + await withWorkspaceTransaction(async (transaction) => { await lockWorkspaceForActor(transaction, this.workspaceId, input.actorUserId); const itemIds = Array.from( new Set(input.relations.flatMap((relation) => [relation.fromItemId, relation.toItemId])), @@ -202,14 +199,7 @@ export class PostgresWorkspacePersistence implements WorkspaceKernelClient { ) .onConflictDoNothing(); } - - const summaries = await getWorkspaceItemsByIds(transaction, this.workspaceId, itemIds); - const itemFacts = await getWorkspaceItemFacts(transaction, this.workspaceId, summaries); - const revision = await nextWorkspaceRevision(transaction, this.workspaceId); - return { result: itemFacts, revision }; }); - await this.notify(command.revision); - return command; } async createItem( @@ -293,7 +283,9 @@ export class PostgresWorkspacePersistence implements WorkspaceKernelClient { }, }; }); - if (outcome.status === "applied") await this.notify(outcome.command.revision); + if (outcome.status === "applied") { + await this.notifyItemsUpserted([outcome.command.result], outcome.command.revision); + } return outcome; } @@ -341,7 +333,9 @@ export class PostgresWorkspacePersistence implements WorkspaceKernelClient { }, }; }); - if (outcome.status === "applied") await this.notify(outcome.command.revision); + if (outcome.status === "applied") { + await this.notifyItemsUpserted([outcome.command.result], outcome.command.revision); + } return outcome; } @@ -430,7 +424,9 @@ export class PostgresWorkspacePersistence implements WorkspaceKernelClient { }, }; }); - if (outcome.status === "applied") await this.notify(outcome.command.revision); + if (outcome.status === "applied") { + await this.notifyItemsUpserted(outcome.command.result, outcome.command.revision); + } return outcome; } @@ -451,7 +447,7 @@ export class PostgresWorkspacePersistence implements WorkspaceKernelClient { const revision = await nextWorkspaceRevision(transaction, this.workspaceId); return { result: item, revision }; }); - await this.notify(command.revision); + await this.notifyItemsUpserted([command.result], command.revision); return command; } @@ -477,7 +473,10 @@ export class PostgresWorkspacePersistence implements WorkspaceKernelClient { await transaction.delete(workspaceItems).where(inArray(workspaceItems.id, rootIds)); } - const revision = await nextWorkspaceRevision(transaction, this.workspaceId); + const revision = + deleteIds.length > 0 + ? await nextWorkspaceRevision(transaction, this.workspaceId) + : await getWorkspaceRevision(transaction, this.workspaceId); const result = { itemIds: rootIds, deletedItemIds: deleteIds }; return { result, @@ -487,12 +486,21 @@ export class PostgresWorkspacePersistence implements WorkspaceKernelClient { }; }); await Promise.all([ - this.notify(command.revision), - this.onItemsDeleted?.({ - workspaceId: this.workspaceId, - documentItemIds: command.documentItemIds, - fileItemIds: command.fileItemIds, - }), + command.result.deletedItemIds.length > 0 + ? this.onChange?.({ + type: "workspace.items.deleted", + workspaceId: this.workspaceId, + revision: command.revision, + itemIds: command.result.deletedItemIds, + }) + : undefined, + command.documentItemIds.length > 0 || command.fileItemIds.length > 0 + ? this.onItemsDeleted?.({ + workspaceId: this.workspaceId, + documentItemIds: command.documentItemIds, + fileItemIds: command.fileItemIds, + }) + : undefined, ]); return { revision: command.revision, result: command.result }; } @@ -531,7 +539,12 @@ export class PostgresWorkspacePersistence implements WorkspaceKernelClient { return await this.files.readPages(input); } - private async notify(revision: number) { - await this.onChange?.({ workspaceId: this.workspaceId, revision }); + private async notifyItemsUpserted(items: WorkspaceItemSummary[], revision: number) { + await this.onChange?.({ + type: "workspace.items.upserted", + workspaceId: this.workspaceId, + revision, + items, + }); } } diff --git a/src/features/workspaces/persistence/workspace-postgres-support.ts b/src/features/workspaces/persistence/workspace-postgres-support.ts index f201203a4..a12c411e8 100644 --- a/src/features/workspaces/persistence/workspace-postgres-support.ts +++ b/src/features/workspaces/persistence/workspace-postgres-support.ts @@ -1,17 +1,14 @@ -import { and, asc, eq, inArray, isNull, or, sql } from "drizzle-orm"; +import { and, asc, eq, inArray, isNull, sql } from "drizzle-orm"; import { workspaceFileAssets, - workspaceItemPages, workspaceItemExtractions, - workspaceItemRelations, workspaceItems, workspaces, } from "#/db/schema"; import { createDbContext, withDb } from "#/db/server"; import type { JsonValue, - WorkspaceItemFacts, WorkspaceItemSummary, WorkspaceItemType, } from "#/features/workspaces/contracts"; @@ -60,13 +57,23 @@ export async function lockWorkspaceForActor( export async function nextWorkspaceRevision(transaction: Transaction, workspaceId: string) { const [workspace] = await transaction .update(workspaces) - .set({ revision: sql`${workspaces.revision} + 1`, updatedAt: new Date() }) + .set({ revision: sql`${workspaces.revision} + 1` }) .where(eq(workspaces.id, workspaceId)) .returning({ revision: workspaces.revision }); if (!workspace) throw new Error("Workspace not found."); return workspace.revision; } +export async function getWorkspaceRevision(db: QueryExecutor, workspaceId: string) { + const [workspace] = await db + .select({ revision: workspaces.revision }) + .from(workspaces) + .where(eq(workspaces.id, workspaceId)) + .limit(1); + if (!workspace) throw new Error("Workspace not found."); + return workspace.revision; +} + export async function getActiveWorkspaceItemRows(db: QueryExecutor, workspaceId: string) { return await db .select() @@ -79,20 +86,15 @@ export async function getActiveWorkspaceItems(db: QueryExecutor, workspaceId: st return (await getActiveWorkspaceItemRows(db, workspaceId)).map(mapWorkspaceItem); } -export async function readWorkspacePageSnapshot(db: QueryExecutor, workspaceId: string) { - const [workspace] = await db - .select({ revision: workspaces.revision }) - .from(workspaces) - .where(eq(workspaces.id, workspaceId)) - .limit(1); - if (!workspace) throw new Error("Workspace not found."); - +export async function readWorkspacePageSnapshot(db: Transaction, workspaceId: string) { + // A transaction owns one pg.Client, so its queries execute serially. Both + // callers use repeatable-read transactions to keep these statements coherent. + const revision = await getWorkspaceRevision(db, workspaceId); const items = await getActiveWorkspaceItems(db, workspaceId); return { workspaceId, items, - itemFacts: await getWorkspaceItemFacts(db, workspaceId, items), - revision: workspace.revision, + revision, }; } @@ -144,53 +146,6 @@ export async function getWorkspaceItemsByIds( }); } -export async function getWorkspaceItemFacts( - db: QueryExecutor, - workspaceId: string, - items: WorkspaceItemSummary[], -): Promise { - if (items.length === 0) return []; - const itemIds = items.map((item) => item.id); - const [relations, pageCounts] = await Promise.all([ - db - .select({ - fromItemId: workspaceItemRelations.fromItemId, - toItemId: workspaceItemRelations.toItemId, - }) - .from(workspaceItemRelations) - .where( - and( - eq(workspaceItemRelations.workspaceId, workspaceId), - or( - inArray(workspaceItemRelations.fromItemId, itemIds), - inArray(workspaceItemRelations.toItemId, itemIds), - ), - ), - ), - db - .select({ - itemId: workspaceItemPages.itemId, - pageCount: sql`count(*)::int`, - }) - .from(workspaceItemPages) - .where(inArray(workspaceItemPages.itemId, itemIds)) - .groupBy(workspaceItemPages.itemId), - ]); - const relationCounts = new Map(); - for (const relation of relations) { - relationCounts.set(relation.fromItemId, (relationCounts.get(relation.fromItemId) ?? 0) + 1); - if (relation.toItemId !== relation.fromItemId) { - relationCounts.set(relation.toItemId, (relationCounts.get(relation.toItemId) ?? 0) + 1); - } - } - const countsByItem = new Map(pageCounts.map((row) => [row.itemId, row.pageCount])); - return items.map((item) => ({ - itemId: item.id, - ...(countsByItem.get(item.id) ? { pageCount: countsByItem.get(item.id) } : {}), - relationshipCount: relationCounts.get(item.id) ?? 0, - })); -} - export async function assertWorkspaceParentIsValid( db: QueryExecutor, workspaceId: string, diff --git a/src/features/workspaces/query-options.ts b/src/features/workspaces/query-options.ts index 0cadc4f13..01185d896 100644 --- a/src/features/workspaces/query-options.ts +++ b/src/features/workspaces/query-options.ts @@ -1,8 +1,10 @@ -import { queryOptions } from "@tanstack/react-query"; +import { queryOptions, replaceEqualDeep } from "@tanstack/react-query"; import { workspacePageQueryKey, workspacesQueryKey } from "#/features/workspaces/cache"; import { getWorkspacePageFn, listWorkspacesFn } from "#/features/workspaces/server/functions"; +type WorkspacePageQueryResult = Awaited>; + export function workspacesQueryOptions() { return queryOptions({ queryKey: workspacesQueryKey, @@ -14,6 +16,15 @@ export function workspacePageQueryOptions(workspaceId: string) { return queryOptions({ queryKey: workspacePageQueryKey(workspaceId), queryFn: () => getWorkspacePageFn({ data: { workspaceId } }), + structuralSharing: (current, incoming) => { + // TanStack exposes this commit-time boundary as unknown even though the + // query function above owns the result type. + const currentPage = current as WorkspacePageQueryResult | undefined; + const incomingPage = incoming as WorkspacePageQueryResult; + return currentPage && incomingPage && currentPage.revision > incomingPage.revision + ? currentPage + : replaceEqualDeep(currentPage, incomingPage); + }, staleTime: 10_000, }); } diff --git a/src/features/workspaces/realtime/messages.ts b/src/features/workspaces/realtime/messages.ts index 82600c559..a5d7275c0 100644 --- a/src/features/workspaces/realtime/messages.ts +++ b/src/features/workspaces/realtime/messages.ts @@ -1,4 +1,5 @@ import { z } from "zod"; +import { workspaceItemSummarySchema } from "#/features/workspaces/contracts"; const workspacePresenceUserSchema = z.object({ id: z.string(), @@ -7,25 +8,44 @@ const workspacePresenceUserSchema = z.object({ image: z.string().nullable(), }); -const workspaceRealtimeServerMessageSchema = z.discriminatedUnion("type", [ +const workspacePageDeltaSchema = z.discriminatedUnion("type", [ z.object({ - type: z.literal("presence.snapshot"), + type: z.literal("workspace.items.upserted"), workspaceId: z.string(), - users: z.array(workspacePresenceUserSchema), + revision: z.number().int().nonnegative(), + items: z.array(workspaceItemSummarySchema).min(1), + }), + z.object({ + type: z.literal("workspace.items.deleted"), + workspaceId: z.string(), + revision: z.number().int().nonnegative(), + itemIds: z.array(z.string()).min(1), }), +]); + +const workspacePageChangeSchema = z.union([ + workspacePageDeltaSchema, z.object({ - type: z.literal("workspace.changed"), + type: z.literal("workspace.page.refresh"), workspaceId: z.string(), revision: z.number().int().nonnegative(), }), ]); +const workspaceRealtimeServerMessageSchema = z.union([ + z.object({ + type: z.literal("presence.snapshot"), + workspaceId: z.string(), + users: z.array(workspacePresenceUserSchema), + }), + workspacePageChangeSchema, +]); + export type WorkspacePresenceUser = z.infer; -export interface WorkspaceRevision { - workspaceId: string; - revision: number; -} +export type WorkspacePageDelta = z.infer; + +export type WorkspacePageChange = z.infer; export interface WorkspaceCommandResult { result: T; diff --git a/src/features/workspaces/realtime/use-workspace-presence.ts b/src/features/workspaces/realtime/use-workspace-presence.ts index 118a0603e..2fc4bfd0f 100644 --- a/src/features/workspaces/realtime/use-workspace-presence.ts +++ b/src/features/workspaces/realtime/use-workspace-presence.ts @@ -6,7 +6,11 @@ import { workspaceKernelAgentName, workspaceKernelBasePath, } from "#/features/workspaces/agent-routes"; -import { parseWorkspaceRealtimeServerMessage, type WorkspacePresenceUser } from "./messages"; +import { + parseWorkspaceRealtimeServerMessage, + type WorkspacePageDelta, + type WorkspacePresenceUser, +} from "./messages"; type ConnectionStatus = "connecting" | "connected" | "disconnected"; @@ -18,8 +22,8 @@ interface PresenceState { interface UseWorkspaceRealtimeInput { workspaceId: string; - lastSeenRevision?: number; - onWorkspaceChanged?: () => void; + onPageChange?: (change: WorkspacePageDelta) => void; + onDesync?: () => void; } function parseServerMessage(data: unknown) { @@ -44,42 +48,28 @@ function getInitialPresenceState(workspaceId: string): PresenceState { export function useWorkspaceRealtime({ workspaceId, - lastSeenRevision, - onWorkspaceChanged, + onPageChange, + onDesync, }: UseWorkspaceRealtimeInput) { const [presence, setPresence] = useState(() => getInitialPresenceState(workspaceId)); - const cachedRevisionRef = useRef(lastSeenRevision ?? 0); - const revisionWorkspaceRef = useRef(workspaceId); - const onWorkspaceChangedRef = useRef(onWorkspaceChanged); + const onPageChangeRef = useRef(onPageChange); + const onDesyncRef = useRef(onDesync); useEffect(() => { - onWorkspaceChangedRef.current = onWorkspaceChanged; + onPageChangeRef.current = onPageChange; + onDesyncRef.current = onDesync; }); const currentPresence = presence.workspaceId === workspaceId ? presence : getInitialPresenceState(workspaceId); - useEffect(() => { - if (revisionWorkspaceRef.current !== workspaceId) { - revisionWorkspaceRef.current = workspaceId; - cachedRevisionRef.current = lastSeenRevision ?? 0; - return; - } - - if (lastSeenRevision === undefined) { - return; - } - - cachedRevisionRef.current = lastSeenRevision; - }, [lastSeenRevision, workspaceId]); - const handleOpen = useCallback(() => { setPresence((current) => ({ ...current, status: "connected", workspaceId, })); - onWorkspaceChangedRef.current?.(); + onDesyncRef.current?.(); }, [workspaceId]); const handleClose = useCallback(() => { @@ -88,7 +78,6 @@ export function useWorkspaceRealtime({ users: [], workspaceId, }); - onWorkspaceChangedRef.current?.(); }, [workspaceId]); const handleError = useCallback(() => { @@ -106,7 +95,7 @@ export function useWorkspaceRealtime({ return; } - if (message?.type === "presence.snapshot" && message.workspaceId === workspaceId) { + if (message.type === "presence.snapshot" && message.workspaceId === workspaceId) { setPresence((current) => ({ ...current, users: message.users, @@ -114,12 +103,12 @@ export function useWorkspaceRealtime({ })); } - if ( - message?.type === "workspace.changed" && - message.workspaceId === workspaceId && - message.revision > cachedRevisionRef.current - ) { - onWorkspaceChangedRef.current?.(); + if (message.type !== "presence.snapshot" && message.workspaceId === workspaceId) { + if (message.type === "workspace.page.refresh") { + onDesyncRef.current?.(); + } else { + onPageChangeRef.current?.(message); + } } }, [workspaceId], diff --git a/src/features/workspaces/realtime/workspace-room-notifier.ts b/src/features/workspaces/realtime/workspace-room-notifier.ts index a8845b8c1..ab88b7631 100644 --- a/src/features/workspaces/realtime/workspace-room-notifier.ts +++ b/src/features/workspaces/realtime/workspace-room-notifier.ts @@ -1,18 +1,18 @@ import { getAgentByName } from "agents"; import { workspaceKernelAgentName } from "#/features/workspaces/agent-routes"; -import type { WorkspaceRevision } from "#/features/workspaces/realtime/messages"; +import type { WorkspacePageChange } from "#/features/workspaces/realtime/messages"; import { recordOperationalFailure } from "#/integrations/observability/operational-events"; const workspaceCleanupDeliveryAttempts = 3; export async function notifyWorkspaceRoom( env: Cloudflare.Env, - change: WorkspaceRevision, + change: WorkspacePageChange, ): Promise { try { const room = await getWorkspaceRoom(env, change.workspaceId); - await room.publishWorkspaceChange(change); + await room.publishWorkspacePageChange(change); } catch (error) { recordOperationalFailure({ error, diff --git a/src/features/workspaces/server/mutations.test.ts b/src/features/workspaces/server/mutations.test.ts index 46c0cd59c..7142d8b7a 100644 --- a/src/features/workspaces/server/mutations.test.ts +++ b/src/features/workspaces/server/mutations.test.ts @@ -88,6 +88,7 @@ describe("workspace settings mutations", () => { ); expect(mocks.nextWorkspaceRevision).toHaveBeenCalledWith(transaction, workspace.id); expect(mocks.notifyWorkspaceRoom).toHaveBeenCalledWith(mocks.env, { + type: "workspace.page.refresh", workspaceId: workspace.id, revision: 8, }); diff --git a/src/features/workspaces/server/mutations.ts b/src/features/workspaces/server/mutations.ts index 2097e1a5e..c79ac51c7 100644 --- a/src/features/workspaces/server/mutations.ts +++ b/src/features/workspaces/server/mutations.ts @@ -187,7 +187,11 @@ export async function updateWorkspaceForCurrentUser( revision: await nextWorkspaceRevision(transaction, input.workspaceId), }; }); - await notifyWorkspaceRoom(env, { workspaceId: input.workspaceId, revision: update.revision }); + await notifyWorkspaceRoom(env, { + type: "workspace.page.refresh", + workspaceId: input.workspaceId, + revision: update.revision, + }); const workspace = { ...update.updatedWorkspace, diff --git a/src/features/workspaces/server/queries.ts b/src/features/workspaces/server/queries.ts index a3bd2edb3..724e3cad6 100644 --- a/src/features/workspaces/server/queries.ts +++ b/src/features/workspaces/server/queries.ts @@ -65,7 +65,6 @@ export async function getWorkspacePageForUser( return { workspace, items: page.items, - itemFacts: page.itemFacts, revision: page.revision, }; }, diff --git a/src/features/workspaces/use-create-workspace.ts b/src/features/workspaces/use-create-workspace.ts index 31e73f6c5..505ce2355 100644 --- a/src/features/workspaces/use-create-workspace.ts +++ b/src/features/workspaces/use-create-workspace.ts @@ -60,7 +60,6 @@ export function useCreateWorkspaceMutation() { setWorkspacePageCache(queryClient, { workspace: optimisticWorkspace, items: [], - itemFacts: [], revision: 0, }); markWorkspaceCreatedThisSession(id); diff --git a/src/features/workspaces/use-workspace-kernel-items.ts b/src/features/workspaces/use-workspace-kernel-items.ts index 6bfc2d2a9..bb74b71bb 100644 --- a/src/features/workspaces/use-workspace-kernel-items.ts +++ b/src/features/workspaces/use-workspace-kernel-items.ts @@ -1,14 +1,12 @@ import { type QueryClient, useMutation, useQueryClient } from "@tanstack/react-query"; import { useServerFn } from "@tanstack/react-start"; -import { useMemo } from "react"; import { toast } from "sonner"; import { + applyWorkspacePageDeltaToCache, createWorkspaceItemInPageCache, - getWorkspaceItemColorInPageCache, moveWorkspaceItemsInPageCache, removeWorkspaceItemsFromPageCache, - updateWorkspaceItemColorInPageCache, workspacePageQueryKey, } from "#/features/workspaces/cache"; import type { @@ -27,9 +25,6 @@ import { updateWorkspaceItemColorFn, } from "#/features/workspaces/server/functions"; import { getErrorMessage } from "#/lib/error-message"; -import { createKeyedDebouncedLatest } from "#/lib/keyed-debounced-latest"; - -const workspaceItemColorCommitDelayMs = 180; export function useCreateWorkspaceItemMutation() { const createWorkspaceItem = useServerFn(createWorkspaceItemFn); const queryClient = useQueryClient(); @@ -49,7 +44,14 @@ export function useCreateWorkspaceItemMutation() { }); } }, - onSuccess: (_command, input) => refreshWorkspacePage(queryClient, input.workspaceId), + onSuccess: (command, input) => { + applyWorkspacePageDeltaToCache(queryClient, { + type: "workspace.items.upserted", + workspaceId: input.workspaceId, + items: [command.result], + revision: command.revision, + }); + }, onError: (error, preparedInput) => { if (preparedInput.id) { removeWorkspaceItemsFromPageCache(queryClient, preparedInput.workspaceId, [ @@ -80,7 +82,14 @@ export function useRenameWorkspaceItemMutation() { return useMutation({ mutationFn: (input: RenameWorkspaceItemInput) => renameWorkspaceItem({ data: input }), - onSuccess: (_command, input) => refreshWorkspacePage(queryClient, input.workspaceId), + onSuccess: (command, input) => { + applyWorkspacePageDeltaToCache(queryClient, { + type: "workspace.items.upserted", + workspaceId: input.workspaceId, + items: [command.result], + revision: command.revision, + }); + }, onError: (error) => { toast.error(getErrorMessage(error, "Unable to rename workspace item right now.")); }, @@ -100,7 +109,14 @@ export function useMoveWorkspaceItemsMutation() { moveWorkspaceItemsInPageCache(queryClient, input); }, - onSuccess: (_command, input) => refreshWorkspacePage(queryClient, input.workspaceId), + onSuccess: (command, input) => { + applyWorkspacePageDeltaToCache(queryClient, { + type: "workspace.items.upserted", + workspaceId: input.workspaceId, + items: command.result, + revision: command.revision, + }); + }, onError: (_error, input) => refreshWorkspacePage(queryClient, input.workspaceId), }); } @@ -109,36 +125,20 @@ export function useUpdateWorkspaceItemColorMutation() { const updateWorkspaceItemColor = useServerFn(updateWorkspaceItemColorFn); const queryClient = useQueryClient(); - const commitColor = useMemo( - () => - createKeyedDebouncedLatest({ - getKey: getWorkspaceItemColorCommitKey, - wait: workspaceItemColorCommitDelayMs, - onExecute: (input) => { - updateWorkspaceItemColor({ data: input }) - .then(() => refreshWorkspacePage(queryClient, input.workspaceId)) - .catch((error: unknown) => { - if (getWorkspaceItemColorInPageCache(queryClient, input) !== input.color) { - return; - } - - void refreshWorkspacePage(queryClient, input.workspaceId); - toast.error(getErrorMessage(error, "Unable to update item color right now.")); - }); - }, - }), - [queryClient, updateWorkspaceItemColor], - ); - - const mutate = (input: UpdateWorkspaceItemColorInput) => { - void queryClient.cancelQueries({ - queryKey: workspacePageQueryKey(input.workspaceId), - }); - updateWorkspaceItemColorInPageCache(queryClient, input); - commitColor(input); - }; - - return { mutate }; + return useMutation({ + mutationFn: (input: UpdateWorkspaceItemColorInput) => updateWorkspaceItemColor({ data: input }), + onSuccess: (command, input) => { + applyWorkspacePageDeltaToCache(queryClient, { + type: "workspace.items.upserted", + workspaceId: input.workspaceId, + items: [command.result], + revision: command.revision, + }); + }, + onError: (error) => { + toast.error(getErrorMessage(error, "Unable to update item color right now.")); + }, + }); } export function useDeleteWorkspaceItemsMutation() { @@ -174,7 +174,16 @@ export function useDeleteWorkspaceItemsMutation() { removeWorkspaceItemsFromPageCache(queryClient, input.workspaceId, input.itemIds); }, - onSuccess: (_command, input) => refreshWorkspacePage(queryClient, input.workspaceId), + onSuccess: (command, input) => { + if (command.result.deletedItemIds.length > 0) { + applyWorkspacePageDeltaToCache(queryClient, { + type: "workspace.items.deleted", + workspaceId: input.workspaceId, + itemIds: command.result.deletedItemIds, + revision: command.revision, + }); + } + }, onError: (_error, input) => refreshWorkspacePage(queryClient, input.workspaceId), }); } @@ -187,10 +196,6 @@ function getDeleteWorkspaceItemsToastMessage( return `${action} ${itemCount === 1 ? "item" : `${itemCount} items`}${suffix}`; } -function getWorkspaceItemColorCommitKey(input: UpdateWorkspaceItemColorInput) { - return `${input.workspaceId}:${input.itemId}`; -} - function refreshWorkspacePage(queryClient: QueryClient, workspaceId: string) { return queryClient.invalidateQueries({ queryKey: workspacePageQueryKey(workspaceId) }); } diff --git a/src/integrations/posthog/ai-observability.ts b/src/integrations/posthog/ai-observability.ts index 5688bb830..1d37cd0f4 100644 --- a/src/integrations/posthog/ai-observability.ts +++ b/src/integrations/posthog/ai-observability.ts @@ -69,21 +69,32 @@ function appendAiTraceProperties( * route that 400'd on every Google leg and was served by OpenAI looked like * Gemini for 105 generations. `credential_type` distinguishes our BYOK keys from * Vercel's metered credits. + * + * `serviceTier` is the tier the provider actually served, not the `priority` we + * asked for: the gateway omits it entirely on a silent downgrade to standard. + * That absence is the signal worth having — priority bills ~1.8-2x when granted, + * so it tells us whether we bought latency or just asked for it. */ export function getGatewayServedRoute(providerMetadata: unknown) { - const routing = (providerMetadata as { gateway?: { routing?: unknown } } | undefined)?.gateway - ?.routing as + const gateway = (providerMetadata as { gateway?: unknown } | undefined)?.gateway as | { - finalProvider?: string; - modelAttempts?: { providerAttempts?: { credentialType?: string; success?: boolean }[] }[]; + serviceTier?: string; + routing?: { + finalProvider?: string; + modelAttempts?: { + providerAttempts?: { credentialType?: string; success?: boolean }[]; + }[]; + }; } | undefined; + const routing = gateway?.routing; return { provider: routing?.finalProvider, credentialType: routing?.modelAttempts ?.flatMap((attempt) => attempt.providerAttempts ?? []) .find((attempt) => attempt.success)?.credentialType, + serviceTier: gateway?.serviceTier, }; } @@ -113,6 +124,7 @@ export function capturePostHogAiGeneration( const properties: Record = { ...captureOptions.properties, ...(served.credentialType ? { credential_type: served.credentialType } : {}), + ...(served.serviceTier ? { service_tier: served.serviceTier } : {}), }; appendAiTraceProperties(properties, { diff --git a/src/integrations/posthog/ai-observability.worker.test.ts b/src/integrations/posthog/ai-observability.worker.test.ts index 912740a16..46152ebda 100644 --- a/src/integrations/posthog/ai-observability.worker.test.ts +++ b/src/integrations/posthog/ai-observability.worker.test.ts @@ -38,6 +38,17 @@ describe("getGatewayServedRoute", () => { ).toBe("system"); }); + // A granted priority tier bills ~2x, and the gateway reports a downgrade by + // omitting the field rather than saying "default". + it("reports the tier that was served, and nothing when priority was downgraded", () => { + const served = routing([{ provider: "openai", credentialType: "byok", success: true }]); + + expect( + getGatewayServedRoute({ gateway: { ...served.gateway, serviceTier: "priority" } }), + ).toEqual({ provider: "openai", credentialType: "byok", serviceTier: "priority" }); + expect(getGatewayServedRoute(served).serviceTier).toBeUndefined(); + }); + it("returns nothing for a non-gateway result, so callers keep their own value", () => { expect(getGatewayServedRoute(undefined)).toEqual({ provider: undefined, diff --git a/src/lib/keyed-debounced-latest.ts b/src/lib/keyed-debounced-latest.ts deleted file mode 100644 index be02668e8..000000000 --- a/src/lib/keyed-debounced-latest.ts +++ /dev/null @@ -1,31 +0,0 @@ -import { debounce } from "@tanstack/pacer"; - -export function createKeyedDebouncedLatest({ - getKey, - onExecute, - wait, -}: { - getKey: (input: TInput) => string; - onExecute: (input: TInput) => void; - wait: number; -}) { - const debouncedByKey = new Map void>(); - - return (input: TInput) => { - const key = getKey(input); - let debounced = debouncedByKey.get(key); - - if (!debounced) { - debounced = debounce( - (latestInput: TInput) => { - debouncedByKey.delete(key); - onExecute(latestInput); - }, - { wait }, - ); - debouncedByKey.set(key, debounced); - } - - debounced(input); - }; -}