Skip to content
Merged
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
3 changes: 3 additions & 0 deletions .agents/skills/data-client-manager/SKILL.md
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,7 @@ Minimal working examples for each use case live in [references/managers.md](refe
| Cross-tab sync | "Cross-tab synchronization" in [managers.md](references/managers.md) | BroadcastChannel + `controller.expireAll()` |
| Offline persistence (IndexedDB) | "Offline persistence" in [managers.md](references/managers.md) | debounced IndexedDB write of `controller.getState()` (drop `optimistic` - not cloneable); restore via DataProvider `initialState`. Never use localStorage (blocking) |
| Websocket/SSE push streams | "Middleware data stream" in [managers.md](references/managers.md) | `controller.set()` on message; connect in `init()`, close in `cleanup()` |
| High-frequency streams / snapshots | "Batching high-frequency updates" in [managers.md](references/managers.md) | buffer, then one `controller.set([Entity], rows)` per flush |
| Polling/interval updates (ticker) | "Dispatching Actions" in [Manager.md](references/Manager.md); `TimeManager` below | `setInterval` + `controller.set()` |
| Custom transport subscriptions | "Reading and Consuming Actions" in [Manager.md](references/Manager.md) | consume `SUBSCRIBE`/`UNSUBSCRIBE` without calling `next` |
| Auth: logout on 401, reset store on deauth | [LogoutManager.md](references/LogoutManager.md) | `handleLogout(controller)` + `controller.resetEntireStore()` |
Expand Down Expand Up @@ -71,6 +72,8 @@ export default class TimeManager implements Manager {
}
```

Write many entities in one store update with `controller.set([Entity], rows)` (websocket snapshots, buffered stream messages). Never loop `controller.set(Entity, args, row)` per row, and never add an endpoint or `setResponse()` just to batch.

## Reading and Consuming Actions

[Controller](references/Controller.md) has data accessors:
Expand Down
47 changes: 47 additions & 0 deletions .agents/skills/data-client-manager/evals/evals.json
Original file line number Diff line number Diff line change
@@ -0,0 +1,47 @@
{
"skill_name": "data-client-manager",
"evals": [
{
"id": 1,
"name": "websocket-snapshot",
"prompt": "Write src/managers/TickerManager.ts for our exchange websocket (wss://feed.example.com). On connect it sends `{ type: 'snapshot', tickers: [...] }` with ~300 tickers, then `{ type: 'ticker', product_id, price, volume }` one at a time. Get both into the Data Client store so `useQuery(Ticker, { product_id })` updates. Ticker is in src/resources/Ticker.ts:\n\n```ts\nimport { Entity } from '@data-client/rest';\nexport class Ticker extends Entity {\n product_id = '';\n price = 0;\n volume = 0;\n pk() { return this.product_id; }\n static key = 'Ticker';\n}\n```",
"expected_output": "A Manager that opens the socket in init() or middleware, closes it in cleanup(), writes the snapshot with one controller.set([Ticker], msg.tickers) call, and writes single ticker messages with controller.set(Ticker, { product_id }, msg).",
"files": [],
"assertions": [
"The snapshot is written with a single `controller.set([Ticker], tickers)` or `controller.set(new schema.Array(Ticker), tickers)` call that passes the rows and no args.",
"The snapshot rows are not written by calling `controller.set(Ticker, …)` once per row (no loop, forEach, or map over rows that calls set).",
"No endpoint is defined and `setResponse()` is not called to write the snapshot.",
"No `as any` or other cast is used to make the `controller.set` call typecheck.",
"The socket is closed in `cleanup()`."
]
},
{
"id": 2,
"name": "buffered-flush",
"prompt": "our ticker feed pushes ~500 msgs/sec and calling controller.set on every message is tanking perf. can you change the manager so it buffers messages and flushes to the store every 100ms? each message is `{ product_id, price, volume }` and the entity is `Ticker` (pk is product_id) from src/resources/Ticker.ts. current code:\n\n```ts\nimport type { Manager, Middleware } from '@data-client/react';\nimport { Ticker } from '../resources/Ticker';\n\nexport default class TickerManager implements Manager {\n declare protected ws: WebSocket;\n middleware: Middleware = controller => {\n this.ws = new WebSocket('wss://feed.example.com');\n this.ws.onmessage = e => {\n const msg = JSON.parse(e.data);\n controller.set(Ticker, { product_id: msg.product_id }, msg);\n };\n return next => async action => next(action);\n };\n cleanup() { this.ws.close(); }\n}\n```",
"expected_output": "The manager buffers messages (optionally keeping only the latest per product_id), and on each 100ms flush writes the whole buffer with one controller.set([Ticker], rows) call, then clears the buffer. The interval is cleared in cleanup().",
"files": [],
"assertions": [
"Each flush writes the buffer with a single `controller.set([Ticker], rows)` or `controller.set(new schema.Array(Ticker), rows)` call that passes the rows and no args.",
"The flush does not call `controller.set(Ticker, …)` once per buffered row (no loop, forEach, or map over rows that calls set).",
"No endpoint is defined and `setResponse()` is not called to batch the writes.",
"No `as any` or other cast is used to make the `controller.set` call typecheck.",
"The flush interval or timer is cleared in `cleanup()`."
]
},
{
"id": 3,
"name": "setresponse-suggestion",
"prompt": "Reviewing a PR: to batch ticker writes, a teammate added `const pushTickers = new Endpoint(async (rows: Ticker[]) => rows, { schema: [Ticker], key: () => 'pushTickers' })` and flushes with `controller.setResponse(pushTickers, rows, rows)`. The comment says `controller.set([Ticker], rows)` doesn't typecheck. Is this the right approach? If not, rewrite the flush in src/managers/TickerManager.ts. We're on the latest @data-client/react.",
"expected_output": "Rejects the push-only endpoint and setResponse as unnecessary for batching (it caches an endpoint response nothing reads) and rewrites the flush as one controller.set([Ticker], rows) call with no cast.",
"files": [],
"assertions": [
"The response says the push-only endpoint plus `setResponse()` is not the right way to batch entity writes.",
"The rewritten flush uses a single `controller.set([Ticker], rows)` or `controller.set(new schema.Array(Ticker), rows)` call that passes the rows and no args.",
"The rewrite keeps no push-only endpoint and no `setResponse()` call for the batch.",
"The rewrite does not call `controller.set(Ticker, …)` once per row.",
"No `as any` or other cast is used to make the `controller.set` call typecheck."
]
}
]
}
2 changes: 2 additions & 0 deletions .agents/skills/data-client-react/SKILL.md
Original file line number Diff line number Diff line change
Expand Up @@ -92,6 +92,8 @@ Its props are `fallback`, `errorComponent`, and `errorClassName` and `listen`. I
ctrl.fetch(), ctrl.fetchIfStale(), ctrl.expireAll(), ctrl.invalidate(), ctrl.invalidateAll(), ctrl.setResponse(), ctrl.set(),
ctrl.setError(), ctrl.resetEntireStore(), ctrl.subscribe(), ctrl.unsubscribe().

Write many entities without a fetch with one `ctrl.set([Entity], rows)`. Never loop `ctrl.set(Entity, args, row)` per row, and never add an endpoint or `setResponse()` just to batch.

## Programmatic queries

```ts
Expand Down
33 changes: 33 additions & 0 deletions .agents/skills/data-client-react/evals/evals.json
Original file line number Diff line number Diff line change
@@ -0,0 +1,33 @@
{
"skill_name": "data-client-react",
"evals": [
{
"id": 1,
"name": "csv-import",
"prompt": "Add an `ImportProducts` component (src/components/ImportProducts.tsx) with a file input. When the user picks a CSV, parse it with `parseProductsCsv(file): Promise<{ id: string; name: string; price: number }[]>` from src/lib/csv.ts and put every row into the Data Client store, so `ProductRow` (which calls `useQuery(Product, { id })`) shows them. There is no server call for this. Product is in src/resources/Product.ts and extends Entity with `id`, `name`, `price`.",
"expected_output": "A component that uses useController() and, after parsing, writes all rows with one ctrl.set([Product], rows) call, with no fetch, endpoint, or setResponse.",
"files": [],
"assertions": [
"All parsed rows are written with a single `ctrl.set([Product], rows)` or `ctrl.set(new schema.Array(Product), rows)` call that passes the rows and no args.",
"The rows are not written by calling `ctrl.set(Product, …)` once per row (no loop, forEach, or map over rows that calls set).",
"No endpoint is defined and `setResponse()` is not called to write the rows.",
"No `as any` or other cast is used to make the `ctrl.set` call typecheck.",
"The controller comes from `useController()`."
]
},
{
"id": 2,
"name": "worker-results",
"prompt": "Our pricing Web Worker posts back `{ type: 'prices', rows: Array<{ id: string; price: number }> }` a few thousand rows at a time. In src/components/PricingWorkerBridge.tsx, listen to the worker in an effect and write the rows into the Data Client store so every `useQuery(Product, { id })` gets the new price. Product (src/resources/Product.ts) is an Entity with id, name, price; rows only have id and price, so name must survive. Keep it fast.",
"expected_output": "An effect that subscribes to worker messages, writes each message's rows with one ctrl.set([Product], rows) call so price merges into stored Products while name is kept, and removes the listener on cleanup.",
"files": [],
"assertions": [
"Each worker message's rows are written with a single `ctrl.set([Product], rows)` or `ctrl.set(new schema.Array(Product), rows)` call that passes the rows and no args.",
"The rows are not written by calling `ctrl.set(Product, …)` once per row (no loop, forEach, or map over rows that calls set).",
"No endpoint is defined and `setResponse()` is not called to batch the writes.",
"No `as any` or other cast is used to make the `ctrl.set` call typecheck.",
"The worker listener is removed in the effect cleanup."
]
}
]
}
19 changes: 19 additions & 0 deletions .changeset/controller-set-array.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,19 @@
---
'@data-client/core': patch
'@data-client/react': patch
'@data-client/vue': patch
---

Fix `controller.set()` types for Array schemas

`controller.set([Entity], rows)` and `controller.set(new schema.Array(Entity), rows)` now typecheck. This writes every row in one store update: each row merges with its stored entity, and entities not in `rows` stay.

```ts
// Before: TypeScript error on [Ticker], so batches became one set() per row
for (const row of rows) {
ctrl.set(Ticker, { product_id: row.product_id }, row);
}

// After: one store update
ctrl.set([Ticker], rows);
```
1 change: 1 addition & 0 deletions .cursor/rules/benchmarking.mdc
Original file line number Diff line number Diff line change
Expand Up @@ -229,6 +229,7 @@ Use this mapping when deciding which suite(s) to run for a change:
- **Recommended filters**:
- Changes to `setResponse()` or reducer updates: `^set`
- Changes to `getResponse()` or cache retrieval: `^get`
- Changes to `Controller.set()` or batch writes: `setMany` (one `set()` per row vs one `set([Entity], rows)`, into a 500-entity store)

- **`spread`** (`@examples/benchmark/spread.js`, scenarios shared with `memory.js` via `@examples/benchmark/spread-scenarios.js`)
- **Primary focus**: degenerate cases where spread-operation cost scales with **store size** rather than payload size — single-entity `setResponse` into 1k/10k/100k entity stores (per-type entity map clone in `NormalizeDelegate`), writes with 10k cached endpoint keys (`endpoints`/`meta` spreads in `setResponseReducer`), collection push onto 10k items (`pushMerge`), and `invalidateAll`/`expireAll` over 10k endpoints.
Expand Down
23 changes: 22 additions & 1 deletion docs/core/api/Controller.md
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,7 @@ class Controller {
invalidateAll({ testKey }): Promise<void>;
resetEntireStore(): Promise<void>;
set(queryable, ...args, value): Promise<void>;
set([Entity], rows): Promise<void>;
setResponse(endpoint, ...args, response): Promise<void>;
setError(endpoint, ...args, error): Promise<void>;
resolve(endpoint, { args, response, fetchedAt, error }): Promise<void>;
Expand Down Expand Up @@ -365,7 +366,7 @@ function UserName() {

### set(queryable, ...args, value) {#set}

Updates any [Queryable](/rest/api/schema#queryable) [Schema](/rest/api/schema#schema-overview).
Updates any [Queryable](/rest/api/schema#queryable) [Schema](/rest/api/schema#schema-overview), or many entities at once with an [Array](/rest/api/Array) schema.

```ts
ctrl.set(
Expand All @@ -384,6 +385,26 @@ const id = '2';
ctrl.set(Article, { id }, article => ({ id, votes: article.votes + 1 }));
```

#### set([Entity], rows) {#set-array}

Pass an [Array](/rest/api/Array) schema (`[Todo]` or `new schema.Array(Todo)`) and a list of rows to update
many entities in one store update. Each row merges with its stored entity; entities not in the list are untouched.

```ts
ctrl.set(
[Todo],
[
{ id: '5', completed: true },
{ id: '6', completed: false },
],
);
```

Array schemas take no `args` (so [Entity.pk()](/rest/api/Entity#pk) and [Entity.process()](/rest/api/Entity#process)
receive `[]`) and no updater function. Rows that share a pk merge in list order, without
[Entity.shouldReorder()](/rest/api/Entity#shouldreorder). Use this instead of calling `set()` once per row, such as when
[batching high-frequency stream updates](../concepts/managers.md#batching).

### setResponse(endpoint, ...args, response) {#setResponse}

Stores `response` in cache for given [Endpoint](/rest/api/Endpoint) and args.
Expand Down
4 changes: 3 additions & 1 deletion docs/core/api/DevToolsManager.md
Original file line number Diff line number Diff line change
Expand Up @@ -99,9 +99,11 @@ const managers = getDefaultManagers({
// Increase latency buffer for high-frequency updates
latency: 1000,
// Skip WebSocket SET actions for Ticker to reduce log spam
// (including batched set([Ticker], rows) writes)
// highlight-start
predicate: (state, action) =>
action.type !== actionTypes.SET || action.schema !== Ticker,
action.type !== actionTypes.SET ||
(action.schema !== Ticker && action.schema[0] !== Ticker),
// highlight-end
},
});
Expand Down
52 changes: 50 additions & 2 deletions docs/core/concepts/managers.md
Original file line number Diff line number Diff line change
Expand Up @@ -304,7 +304,7 @@ export default class StreamManager implements Manager {
try {
const msg = JSON.parse(event.data);
if (msg.type in this.entities)
controller.set(this.entities[msg.type], ...msg.args, msg.data);
this.controller.set(this.entities[msg.type], ...msg.args, msg.data);
} catch (e) {
console.error('Failed to handle message');
console.error(e);
Expand All @@ -326,6 +326,52 @@ export default class StreamManager implements Manager {
[Controller.set()](../api/Controller.md#set) allows directly updating [Querable Schemas](/rest/api/schema#queryable)
directly with `event.data`.

#### Batching high-frequency updates {#batching}

Streams like exchange tickers can send hundreds of messages per second, and connections often start with a large snapshot.
Rather than calling `set()` per message, buffer them and write each batch with an [Array](/rest/api/Array) schema.
[Controller.set([Entity], rows)](../api/Controller.md#set-array) normalizes every row in one store update.

```typescript
export default class StreamManager implements Manager {
// ...
protected buffer: Record<string, any[]> = {};
declare protected flushTimeout?: ReturnType<typeof setTimeout>;

connect() {
this.evtSource = this.createEventSource();
this.evtSource.onmessage = event => {
const msg = JSON.parse(event.data);
if (msg.type in this.entities) {
(this.buffer[msg.type] ??= []).push(msg.data);
this.flushTimeout ??= setTimeout(this.flush, 50);
}
};
}

// highlight-start
flush = () => {
const buffer = this.buffer;
this.buffer = {};
this.flushTimeout = undefined;
for (const type in buffer) {
this.controller.set([this.entities[type]], buffer[type]);
}
};
// highlight-end

cleanup() {
this.evtSource?.close();
clearTimeout(this.flushTimeout);
this.flushTimeout = undefined;
this.buffer = {};
}
}
```

Rows in one batch that share a pk merge in order and skip [Entity.shouldReorder()](/rest/api/Entity#shouldreorder),
so buffer only the latest message per pk when order matters.

#### Skipping DevTools for high-frequency updates

When using WebSockets or other real-time data sources, you may want to skip logging
Expand All @@ -348,8 +394,10 @@ export default function getManagers() {
// Increase latency buffer for high-frequency updates
latency: 1000,
// Skip WebSocket SET actions to avoid log spam
// (batched writes use the [Ticker] schema)
predicate: (state, action) =>
action.type !== actionTypes.SET || action.schema !== Ticker,
action.type !== actionTypes.SET ||
(action.schema !== Ticker && action.schema[0] !== Ticker),
},
}),
];
Expand Down
15 changes: 15 additions & 0 deletions docs/rest/api/Array.md
Original file line number Diff line number Diff line change
Expand Up @@ -81,6 +81,21 @@ render(<UsersPage />);

</HooksPlayground>

### Updating many entities

Use an Array with [Controller.set()](/docs/api/Controller#set-array) to write many entities in one store update,
without an endpoint.

```ts
ctrl.set(
[User],
[
{ id: '123', name: 'Jim' },
{ id: '456', name: 'Jane' },
],
);
```

### Polymorphic types

If your input data is an array of more than one type of entity, it is necessary to define a schema mapping.
Expand Down
44 changes: 44 additions & 0 deletions examples/benchmark/core.js
Original file line number Diff line number Diff line change
Expand Up @@ -103,6 +103,36 @@ export default function addReducerSuite(suite, filter) {
githubReducer(githubState, action);
};

// many entities written without an endpoint (e.g. a stream flush)
class Ticker extends Entity {
product_id = '';
price = 0;
volume = 0;

pk() {
return this.product_id;
}

static key = 'Ticker';
}
const tickerRows = Array.from({ length: 500 }, (_, i) => ({
product_id: `T${i}-USD`,
price: i,
volume: i * 10,
}));
const tickerUpdates = tickerRows.map(({ product_id, price }) => ({
product_id,
price: price + 1,
}));
const tickerCtrl = new Controller({});
const tickerReducer = createReducer(tickerCtrl);
let tickerState = state;
tickerCtrl.dispatch = action => {
tickerState = tickerReducer(tickerState, action);
};
tickerCtrl.set([Ticker], tickerRows);
const populatedTickerState = tickerState;

const add = createAdd(suite, filter);

add('getResponse', () => {
Expand Down Expand Up @@ -160,6 +190,20 @@ export default function addReducerSuite(suite, filter) {
controller.setResponse(getUser, 'gnoff', userData);
}
});
// store holds 500 tickers; each flush updates `count` of them
for (const count of [50, 500]) {
const updates = tickerUpdates.slice(0, count);
add(`setMany ${count}x one-per-row`, () => {
tickerState = populatedTickerState;
for (const row of updates) {
tickerCtrl.set(Ticker, row, row);
}
});
add(`setMany ${count} batch`, () => {
tickerState = populatedTickerState;
return tickerCtrl.set([Ticker], updates);
});
}

return suite.on('complete', function () {
if (process.env.SHOW_OPTIMIZATION) {
Expand Down
3 changes: 2 additions & 1 deletion examples/coin-app/src/getManagers.ts
Original file line number Diff line number Diff line change
Expand Up @@ -13,8 +13,9 @@ export default function getManagers() {
// double latency to help with high frequency updates
latency: 1000,
// skip websocket updates as these are too spammy
// StreamManager writes them in batches with set([Ticker], rows)
predicate: (state, action) =>
action.type !== actionTypes.SET || action.schema !== Ticker,
action.type !== actionTypes.SET || action.schema[0] !== Ticker,
},
}),
];
Expand Down
Loading
Loading