Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 2 additions & 2 deletions packages/devframe/src/client/rpc-shared-state.ts
Original file line number Diff line number Diff line change
Expand Up @@ -121,7 +121,7 @@ export function createRpcSharedStateClientHost(rpc: DevframeRpcClient): RpcShare
}
}

return new Promise<SharedState<T>>((resolve) => {
return new Promise<SharedState<T>>((resolve, reject) => {
if (!rpc.isTrusted) {
resolve(state)
let initialized = false
Expand All @@ -133,7 +133,7 @@ export function createRpcSharedStateClientHost(rpc: DevframeRpcClient): RpcShare
})
}
else {
initSharedState().then(resolve)
initSharedState().then(resolve, reject)
}
})
},
Expand Down
3 changes: 2 additions & 1 deletion packages/devframe/src/client/scope.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@ function createMockClient() {

// eslint-disable-next-line slop/no-chained-type-assertions -- partial test mock exercises only the members client.scope() touches
const rpc = {
ensureTrusted: vi.fn(async () => true),
call: vi.fn((..._args: any[]) => Promise.resolve('ok')),
callEvent: vi.fn((..._args: any[]) => {}),
callOptional: vi.fn((..._args: any[]) => Promise.resolve('ok')),
Expand Down Expand Up @@ -126,7 +127,7 @@ describe('client.scope()', () => {
const { settings } = createScopedClientContext(rpc, 'my-plugin')

await settings.global.set('token', 'abc')
expect(rpc.sharedState.get).toHaveBeenCalledWith('devframe:settings:global:my-plugin', { initialValue: {} })
expect(rpc.sharedState.get).toHaveBeenCalledWith('devframe:settings:global:my-plugin')
expect(await settings.global.get('token')).toBe('abc')

await settings.project.set('theme', 'dark')
Expand Down
153 changes: 153 additions & 0 deletions packages/devframe/src/client/settings.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,153 @@
import type { DevframeRpcClient } from './rpc'
import { existsSync, mkdtempSync, readFileSync, rmSync } from 'node:fs'
import { tmpdir } from 'node:os'
import { join } from 'node:path'
import { setImmediate } from 'node:timers/promises'
import { createEventEmitter } from 'devframe/utils/events'
import { describe, expect, it, vi } from 'vitest'
import { DEVFRAME_EVENTS } from '../events'
import { createHostContext } from '../node/context'
import { createRpcSharedStateClientHost } from './rpc-shared-state'
import { createClientSettings } from './settings'

function setup(backend = 'live', trusted = true) {
const snapshot = Promise.withResolvers<Record<string, unknown> | undefined>()
const trust = Promise.withResolvers<boolean>()
const handlers = new Map<string, (...args: any[]) => void>()
// eslint-disable-next-line slop/no-chained-type-assertions -- partial RPC mock retains the shared-state implementation
const rpc = {
connectionMeta: { backend },
isTrusted: trusted,
events: createEventEmitter<any>(),
client: { register: (fn: any) => handlers.set(fn.name, fn.handler) },
ensureTrusted: vi.fn(() => trusted ? Promise.resolve(true) : trust.promise),
call: vi.fn(() => snapshot.promise),
callEvent: vi.fn(),
} as unknown as DevframeRpcClient
rpc.sharedState = createRpcSharedStateClientHost(rpc)
const settings = createClientSettings(rpc, 'test')
return { rpc, settings, snapshot, trust, handlers }
}

describe('client settings initialization', () => {
it.each(['project', 'global'] as const)('persists the first %s write after a delayed empty snapshot', async (scope) => {
const dir = mkdtempSync(join(tmpdir(), 'devframe-settings-'))
const host = {
mountStatic: () => {},
resolveOrigin: () => 'http://localhost',
getStorageDir: (storageScope: string) => join(dir, storageScope),
}
const ctx = await createHostContext({ cwd: dir, mode: 'dev', host })
const snapshotReady = Promise.withResolvers<void>()
const releaseSnapshot = Promise.withResolvers<void>()
const writes: Promise<unknown>[] = []
// eslint-disable-next-line slop/no-chained-type-assertions -- control transport timing while retaining both shared-state implementations
const rpc = {
connectionMeta: { backend: 'live' },
isTrusted: true,
ensureTrusted: async () => true,
client: { register: () => {} },
call: async (method: string, key: string) => {
const value = await ctx.rpc.invokeLocal(method as any, key)
snapshotReady.resolve()
await releaseSnapshot.promise
return value
},
callEvent: (method: string, ...args: any[]) => {
if (method !== 'devframe:rpc:server-state:subscribe')
writes.push(ctx.rpc.invokeLocal(method as any, ...args))
},
} as unknown as DevframeRpcClient
rpc.sharedState = createRpcSharedStateClientHost(rpc)
const settings = createClientSettings(rpc, 'new-tool')[scope]

try {
const write = settings.set('theme', 'dark')
await snapshotReady.promise
await setImmediate()
releaseSnapshot.resolve()
await write
await setImmediate()
await Promise.all(writes)

const filepath = join(dir, scope, 'settings/new-tool.json')
await expect.poll(() => existsSync(filepath)).toBe(true)
expect(await settings.all()).toEqual({ theme: 'dark' })
expect(JSON.parse(readFileSync(filepath, 'utf8'))).toEqual({ theme: 'dark' })
const restarted = await createHostContext({ cwd: dir, mode: 'dev', host })
expect(await restarted.scope('new-tool').settings[scope].all()).toEqual({ theme: 'dark' })
}
finally {
rmSync(dir, { recursive: true, force: true })
}
})

it.each(['project', 'global'] as const)('reads the initial %s snapshot before resolving', async (scope) => {
const { settings, snapshot } = setup()
const read = settings[scope].get('theme')
await setImmediate()
snapshot.resolve({ theme: 'light' })
await expect(read).resolves.toBe('light')
})

it.each(['project', 'global'] as const)('preserves an immediate %s write and unrelated settings', async (scope) => {
const { rpc, settings, snapshot } = setup()
const write = settings[scope].set('theme', 'dark')
const read = settings[scope].all()
await setImmediate()
snapshot.resolve({ theme: 'light', language: 'en' })
await write
await expect(settings[scope].get('theme')).resolves.toBe('dark')
await expect(read).resolves.toEqual({ theme: 'dark', language: 'en' })
expect(rpc.call).toHaveBeenCalledTimes(1)
})

it('applies an immediate deletion after the initial snapshot', async () => {
const { settings, snapshot } = setup()
const deletion = settings.project.delete('theme')
await setImmediate()
snapshot.resolve({ theme: 'light', language: 'en' })
await deletion
await expect(settings.project.all()).resolves.toEqual({ language: 'en' })
})

it.each(['live', 'static'])('starts an absent %s store with an empty object', async (backend) => {
const { settings, snapshot } = setup(backend)
const read = settings.project.all()
await setImmediate()
snapshot.resolve(undefined)
await expect(read).resolves.toEqual({})
await settings.project.set('theme', 'dark')
await expect(settings.project.get('theme')).resolves.toBe('dark')
})

it('waits for trust before requesting the initial snapshot', async () => {
const { rpc, settings, snapshot, trust } = setup('live', false)
const write = settings.project.set('theme', 'dark')
await setImmediate()
expect(rpc.call).not.toHaveBeenCalled()
expect(rpc.callEvent).not.toHaveBeenCalled()
Object.defineProperty(rpc, 'isTrusted', { value: true })
trust.resolve(true)
snapshot.resolve({ language: 'en' })
await write
await expect(settings.project.all()).resolves.toEqual({ theme: 'dark', language: 'en' })
})

it('rejects a failed initial snapshot and retries on the next operation', async () => {
const { rpc, settings } = setup()
vi.mocked(rpc.call).mockRejectedValueOnce(new Error('snapshot failed'))
await expect(settings.project.get('theme')).rejects.toThrow('snapshot failed')
vi.mocked(rpc.call).mockResolvedValueOnce({ theme: 'light' })
await expect(settings.project.get('theme')).resolves.toBe('light')
})

it('continues to accept remote updates after initialization', async () => {
const { settings, snapshot, handlers } = setup()
await setImmediate()
snapshot.resolve({ theme: 'light' })
await settings.project.set('theme', 'dark')
handlers.get(DEVFRAME_EVENTS.broadcast.clientStateUpdated)!('devframe:settings:project:test', { theme: 'system' }, 'remote')
await expect(settings.project.get('theme')).resolves.toBe('system')
})
})
17 changes: 12 additions & 5 deletions packages/devframe/src/client/settings.ts
Original file line number Diff line number Diff line change
Expand Up @@ -11,13 +11,20 @@ function createClientSettingsStore<T extends Record<string, any>>(
const stateKey = `devframe:settings:${scope}:${namespace}`
let statePromise: Promise<SharedState<T>> | undefined

// The client mirrors the server's file-backed settings store over the
// shared-state sync protocol: providing an empty initial value lets the
// client subscribe and merge the authoritative server snapshot, and any
// local `set` is pushed back to (and persisted by) the server.
// Wait for the initial snapshot before reading or changing settings.
// This keeps initialization from overwriting the first local operation.
function store(): Promise<SharedState<T>> {
if (!statePromise) {
statePromise = (rpc.sharedState.get as any)(stateKey, { initialValue: {} }) as Promise<SharedState<T>>
statePromise = (async () => {
await rpc.ensureTrusted()
const state = await rpc.sharedState.get<T>(stateKey)
if (state.value() === undefined)
state.mutate(() => ({} as T))
return state
})().catch((error) => {
statePromise = undefined
throw error
})
}
return statePromise
}
Expand Down
36 changes: 36 additions & 0 deletions packages/devframe/src/node/__tests__/scope.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -180,6 +180,42 @@ describe('ctx.scope()', () => {
expect(JSON.parse(readFileSync(globalFile, 'utf-8'))).toEqual({ token: 'abc' })
})

it.each(['project', 'global'] as const)('persists client-first %s settings and loads them on first RPC read', async (scope) => {
const { ctx, dir } = await createCtx()
const key = `devframe:settings:${scope}:client-only`
// Exercise the handlers used by RPC clients without touching node settings first.
await ctx.rpc.invokeLocal('devframe:rpc:server-state:set', key, { theme: 'dark' }, 'client')
await sleep(250)
const restarted = await createHostContext({ cwd: dir, mode: 'dev', host: createTestHost(dir) })
expect(await restarted.rpc.invokeLocal('devframe:rpc:server-state:get', key)).toEqual({ theme: 'dark' })
expect(await restarted.scope('client-only').settings[scope].all()).toEqual({ theme: 'dark' })
})

it('loads existing settings for a client-first patch and shares the same state with node settings', async () => {
const { ctx, dir } = await createCtx()
await ctx.scope('my-plugin').settings.project.set('theme', 'dark')
await sleep(250)
const restarted = await createHostContext({ cwd: dir, mode: 'dev', host: createTestHost(dir) })
const key = 'devframe:settings:project:my-plugin'
await restarted.rpc.invokeLocal('devframe:rpc:server-state:patch', key, [{ op: 'add', path: ['zoom'], value: 2 }], 'client')
expect(await restarted.scope('my-plugin').settings.project.all()).toEqual({ theme: 'dark', zoom: 2 })
expect(await restarted.scope('other-plugin').settings.project.all()).toEqual({})
expect(await restarted.scope('my-plugin').settings.global.all()).toEqual({})
await sleep(250)
expect(JSON.parse(readFileSync(join(dir, 'project/settings/my-plugin.json'), 'utf-8'))).toEqual({ theme: 'dark', zoom: 2 })
})

it.each(['ordinary:state', 'devframe:settings:workspace:plugin', 'devframe:settings:project:../escaped'])('keeps %s in memory', async (key) => {
const { ctx, dir } = await createCtx()
expect(await ctx.rpc.invokeLocal('devframe:rpc:server-state:get', key)).toBeUndefined()
await ctx.rpc.invokeLocal('devframe:rpc:server-state:set', key, { value: 1 }, 'client')
expect(await ctx.rpc.invokeLocal('devframe:rpc:server-state:get', key)).toEqual({ value: 1 })
await sleep(250)
expect(existsSync(join(dir, 'project'))).toBe(false)
const restarted = await createHostContext({ cwd: dir, mode: 'dev', host: createTestHost(dir) })
expect(await restarted.rpc.invokeLocal('devframe:rpc:server-state:get', key)).toBeUndefined()
})

it('notifies onChange subscribers', async () => {
const { ctx } = await createCtx()
const { settings } = ctx.scope('my-plugin')
Expand Down
3 changes: 2 additions & 1 deletion packages/devframe/src/node/host-functions.ts
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@ import { removeClientAgentSession } from './client-agent'
import { diagnostics } from './diagnostics'
import { createRpcSharedStateServerHost } from './rpc-shared-state'
import { createRpcStreamingServerHost } from './rpc-streaming'
import { resolveSettingsState } from './settings'

const debugBroadcast = createDebug('devframe:rpc:broadcast')

Expand Down Expand Up @@ -38,7 +39,7 @@ export class RpcFunctionsHostImpl extends RpcFunctionsCollectorBase<DevframeRpcS
constructor(context: DevframeNodeContext) {
super(context)

this.sharedState = createRpcSharedStateServerHost(this)
this.sharedState = createRpcSharedStateServerHost(this, key => resolveSettingsState(context, key))
this.streaming = createRpcStreamingServerHost(this)
}

Expand Down
44 changes: 25 additions & 19 deletions packages/devframe/src/node/rpc-shared-state.ts
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@ const debugSubscribe = createDebug('devframe:rpc:state:subscribe')

export function createRpcSharedStateServerHost(
rpc: RpcFunctionsHost,
resolveState?: (key: string) => SharedState<any> | undefined,
): RpcSharedStateHost {
const sharedState = new Map<string, SharedState<any>>()
const stateDisposers = new Map<string, () => void>()
Expand Down Expand Up @@ -46,24 +47,35 @@ export function createRpcSharedStateServerHost(
}
}

function addState(key: string, state: SharedState<any>) {
debug('new-state', key)
stateDisposers.set(key, registerSharedState(key, state))
sharedState.set(key, state)
for (const fn of keyAddedListeners)
fn(key)
return state
}

function resolve(key: string) {
const existing = sharedState.get(key)
if (existing)
return existing
const state = resolveState?.(key)
return state ? addState(key, state) : undefined
}

const host: RpcSharedStateHost = {
get: async <T extends object>(key: string, options?: RpcSharedStateGetOptions<T>) => {
if (sharedState.has(key)) {
return sharedState.get(key)!
}
const existing = resolve(key)
if (existing)
return existing
if (options?.initialValue === undefined && options?.sharedState === undefined) {
throw diagnostics.DF0013({ key })
}
debug('new-state', key)
const state = options.sharedState ?? createSharedState<T>({
return addState(key, options.sharedState ?? createSharedState<T>({
initialValue: options.initialValue as T,
enablePatches: false,
})
stateDisposers.set(key, registerSharedState(key, state))
sharedState.set(key, state)
for (const fn of keyAddedListeners)
fn(key)
return state
}))
},
keys() {
return Array.from(sharedState.keys())
Expand Down Expand Up @@ -106,10 +118,7 @@ export function createRpcSharedStateServerHost(
name: 'devframe:rpc:server-state:get',
type: 'query',
handler: async (key: string) => {
if (!sharedState.has(key))
return undefined
const state = await host.get(key)
return state.value()
return resolve(key)?.value()
},
/**
* Pre-compute snapshots for the build-mode static dump so the SPA
Expand Down Expand Up @@ -139,10 +148,7 @@ export function createRpcSharedStateServerHost(
name: 'devframe:rpc:server-state:patch',
type: 'query',
handler: async (key: string, patches: SharedStatePatch[], syncId: string) => {
if (!sharedState.has(key))
return
const state = await host.get(key)
state.patch(patches, syncId)
resolve(key)?.patch(patches, syncId)
},
})

Expand Down
Loading
Loading