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
5 changes: 5 additions & 0 deletions .changeset/reject-reused-sync-rows.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
---
'@tanstack/db': patch
---

Throw `SyncRowReusedWithoutPreviousValueError` in development when a sync source changes a top-level field of a row object it already wrote and writes it again without `previousValue`. The collection keeps the written object as the stored row, so the in-place change overwrote the previous value, and live queries could keep the row in a filter it left. The check compares shallow copies, so it does not detect a change inside a nested object. Production builds skip the check.
1 change: 1 addition & 0 deletions docs/contributing/oracle-coverage.md

Large diffs are not rendered by default.

17 changes: 17 additions & 0 deletions docs/guides/collection-options-creator.md
Original file line number Diff line number Diff line change
Expand Up @@ -180,6 +180,23 @@ applying its rows. A successful `loadSubset` must await or return every commit
receipt that establishes its result. Do not use `begin({ immediate: true })` to
bypass that ordering just to settle a load.

The collection keeps the object you pass to `write()` as the row's stored
value. Pass a new object for each update. If your source changes a row object
in place and writes it again, pass the row's previous value as
`previousValue`. Without it, the change already overwrote the value the
collection needs to publish the update, and live queries can keep the row in
a result it left. In development, the collection throws
`SyncRowReusedWithoutPreviousValueError` for that write when a top-level field
changed. The check compares shallow copies, so it does not detect a change
inside a nested object. Pass `previousValue` for those writes too.

```ts
// Changes a stored row in place, so it must name the previous value
const previousValue = { ...row }
row.status = `done`
write({ type: `update`, value: row, previousValue })
```

For request-scoped writes, pass the request's abort signal to `commit(signal)`.
Cancellation before application rejects the receipt with `AbortError`; aborting
after application does not undo published rows. Do not attach one request's
Expand Down
5 changes: 4 additions & 1 deletion packages/db/mangle-cache.json
Original file line number Diff line number Diff line change
Expand Up @@ -366,5 +366,8 @@
"canRetryRepair": "f0",
"restore": "f1",
"attach": "f2",
"compilations": "f3"
"compilations": "f3",
"writtenRows": "f4",
"equalityRoute": "f5",
"checkReusedRow": "f6"
}
60 changes: 60 additions & 0 deletions packages/db/src/collection/sync.ts
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@ import {
NoPendingSyncTransactionCommitError,
NoPendingSyncTransactionWriteError,
SyncCleanupError,
SyncRowReusedWithoutPreviousValueError,
SyncTransactionAlreadyCommittedError,
SyncTransactionAlreadyCommittedWriteError,
} from '../errors'
Expand Down Expand Up @@ -46,6 +47,27 @@ type LoadSubsetOperation = {
deferred?: Deferred<void>
}

function shallowEqual(
left: Record<string, unknown>,
right: Record<string, unknown>,
): boolean {
const keys = Object.keys(left)
return (
keys.length === Object.keys(right).length &&
keys.every((key) => Object.is(left[key], right[key]))
)
}

// Bundlers inline `process.env.NODE_ENV`; without a bundler or `process`,
// the development checks stay off.
function isDevelopment(): boolean {
try {
return process.env.NODE_ENV !== `production`
} catch {
return false
}
}

export class CollectionSyncManager<
TOutput extends object = Record<string, unknown>,
TKey extends string | number = string | number,
Expand All @@ -59,6 +81,8 @@ export class CollectionSyncManager<
private config!: CollectionConfig<TOutput, TKey, TSchema, any>
private id: string
private syncMode: `eager` | `on-demand`
// Development only: each written row object's fields when it was written.
private writtenRows: WeakMap<object, Record<string, unknown>> | undefined

public preloadPromise: Promise<void> | null = null
private rejectPreload?: (error: unknown) => void
Expand Down Expand Up @@ -111,6 +135,38 @@ export class CollectionSyncManager<
this._events = deps.events
}

/**
* Core keeps the object a source writes as the stored row. A source that
* changes that object in place and writes it again has already
* overwritten the previous value core would publish, unless it passes
* `previousValue`. Rewriting an unchanged object stays valid.
*/
private checkReusedRow(
key: TKey,
type: string,
message: { value: TOutput; previousValue?: TOutput },
): void {
const value = message.value as Record<string, unknown>
const writtenRows = (this.writtenRows ??= new WeakMap())
try {
const written = writtenRows.get(value)
if (
written &&
type === `update` &&
// A write that names its previous value declares the reuse.
!(`previousValue` in message) &&
!shallowEqual(written, value)
) {
throw new SyncRowReusedWithoutPreviousValueError(key)
}
writtenRows.set(value, { ...value })
} catch (error) {
if (error instanceof SyncRowReusedWithoutPreviousValueError) throw error
// A row that cannot be read reports its failure where it is applied.
writtenRows.delete(value)
}
}

private createDuplicateKeyError(key: TKey): DuplicateKeySyncError {
const utils = this.config.utils as
| Partial<LiveQueryCollectionUtils>
Expand Down Expand Up @@ -219,6 +275,10 @@ export class CollectionSyncManager<
messageType = disposition
}

if (`value` in messageWithOptionalKey && isDevelopment()) {
this.checkReusedRow(key, messageType, messageWithOptionalKey)
}

const message = {
...messageWithOptionalKey,
type: messageType,
Expand Down
10 changes: 10 additions & 0 deletions packages/db/src/errors.ts
Original file line number Diff line number Diff line change
Expand Up @@ -393,6 +393,16 @@ export class SyncTransactionAlreadyCommittedWriteError extends TransactionError
}
}

export class SyncRowReusedWithoutPreviousValueError extends TransactionError {
constructor(key: string | number) {
super(
`A sync update for key "${key}" wrote a row object that changed in place since it was last written. ` +
`The change overwrote the row's previous value. ` +
`Write a new object, or pass the row's previous value as \`previousValue\`.`,
)
}
}

export class NoPendingSyncTransactionCommitError extends TransactionError {
constructor() {
super(`No pending sync transaction to commit`)
Expand Down
110 changes: 110 additions & 0 deletions packages/db/tests/sync-reused-row.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,110 @@
import { afterEach, describe, expect, it, vi } from 'vitest'
import { createCollection } from '../src/collection/index.js'
import { SyncRowReusedWithoutPreviousValueError } from '../src/errors.js'
import { createLiveQueryCollection, eq } from '../src/query/index.js'
import type { SyncConfig } from '../src/types.js'

/**
* A sync update tells the Collection a row's new value. Core keeps the
* object a source writes as the row's stored value, so a source that
* changes that object in place and writes it again has already overwritten
* the previous value core would publish. Live queries then see an update
* whose old and new values are the same object, and a row that left a
* filter stays in it. A source that reuses its row object must say what the
* row was through `previousValue`. In development, core rejects the write
* that omits it; production keeps the cheaper unchecked path.
*/
type Row = { id: string; group: string }

function setup() {
let sync!: Parameters<SyncConfig<Row, string>[`sync`]>[0]
const row: Row = { id: `r`, group: `a` }
const collection = createCollection<Row, string>({
id: `sync-reused-row`,
getKey: (item) => item.id,
startSync: true,
sync: {
rowUpdateMode: `full`,
sync: (actions) => {
sync = actions
actions.begin()
actions.write({ type: `insert`, value: row })
actions.commit()
actions.markReady()
},
},
})
const groupA = createLiveQueryCollection({
query: (q) => q.from({ r: collection }).where(({ r }) => eq(r.group, `a`)),
startSync: true,
})
return { sync: () => sync, row, collection, groupA }
}

describe(`sync writes of a reused row object`, () => {
afterEach(() => vi.unstubAllEnvs())

it(`rejects an in-place update without previousValue in development`, async () => {
vi.stubEnv(`NODE_ENV`, `development`)
const { sync, row, collection, groupA } = setup()
await groupA.preload()
row.group = `b`
sync().begin()
expect(() => sync().write({ type: `update`, value: row })).toThrow(
SyncRowReusedWithoutPreviousValueError,
)
await collection.cleanup()
})

it(`accepts an in-place update that names its previous value`, async () => {
vi.stubEnv(`NODE_ENV`, `development`)
const { sync, row, collection, groupA } = setup()
await groupA.preload()
expect([...groupA.keys()]).toEqual([`r`])
const previousValue = { ...row }
row.group = `b`
sync().begin()
sync().write({ type: `update`, value: row, previousValue })
sync().commit()
expect([...groupA.keys()]).toEqual([])
await collection.cleanup()
})

it(`accepts an unchanged rewrite after a declared in-place update`, async () => {
vi.stubEnv(`NODE_ENV`, `development`)
const { sync, row, collection, groupA } = setup()
await groupA.preload()
const previousValue = { ...row }
row.group = `b`
sync().begin()
sync().write({ type: `update`, value: row, previousValue })
sync().commit()
// The live-query Collection rewrites an unchanged row the same way.
sync().begin()
expect(() => sync().write({ type: `update`, value: row })).not.toThrow()
sync().commit()
await collection.cleanup()
})

it(`accepts an update with a new object`, async () => {
vi.stubEnv(`NODE_ENV`, `development`)
const { sync, collection, groupA } = setup()
await groupA.preload()
sync().begin()
sync().write({ type: `update`, value: { id: `r`, group: `b` } })
sync().commit()
expect([...groupA.keys()]).toEqual([])
await collection.cleanup()
})

it(`does not check in production`, async () => {
vi.stubEnv(`NODE_ENV`, `production`)
const { sync, row, collection, groupA } = setup()
await groupA.preload()
row.group = `b`
sync().begin()
expect(() => sync().write({ type: `update`, value: row })).not.toThrow()
sync().commit()
await collection.cleanup()
})
})
2 changes: 1 addition & 1 deletion scripts/test-minified-db.mjs
Original file line number Diff line number Diff line change
Expand Up @@ -84,7 +84,7 @@ const errorGroups = {
CollectionStateError: `CollectionStateError CollectionInErrorStateError InvalidCollectionStatusTransitionError CollectionIsInErrorStateError NegativeActiveSubscribersError LiveQueryObserverDisposedError LiveQueryWindowControllerDisposedError`,
CollectionOperationError: `CollectionOperationError UndefinedKeyError InvalidKeyError DuplicateKeyError DuplicateKeySyncError MissingUpdateArgumentError NoKeysPassedToUpdateError UpdateKeyNotFoundError KeyUpdateNotAllowedError NoKeysPassedToDeleteError DeleteKeyNotFoundError`,
MissingHandlerError: `MissingHandlerError MissingInsertHandlerError MissingUpdateHandlerError MissingDeleteHandlerError`,
TransactionError: `TransactionError MissingMutationFunctionError TransactionNotPendingMutateError TransactionAlreadyCompletedRollbackError TransactionNotPendingCommitError NoPendingSyncTransactionWriteError SyncTransactionAlreadyCommittedWriteError NoPendingSyncTransactionCommitError SyncTransactionAlreadyCommittedError`,
TransactionError: `TransactionError MissingMutationFunctionError TransactionNotPendingMutateError TransactionAlreadyCompletedRollbackError TransactionNotPendingCommitError NoPendingSyncTransactionWriteError SyncTransactionAlreadyCommittedWriteError SyncRowReusedWithoutPreviousValueError NoPendingSyncTransactionCommitError SyncTransactionAlreadyCommittedError`,
QueueCapacityExceededError: `QueueCapacityExceededError`,
QueueDisposedError: `QueueDisposedError`,
ThrottleCallDroppedError: `ThrottleCallDroppedError`,
Expand Down
Loading