Skip to content

perf(agents-server): resolve entity schemas once per stream append - #4821

Open
adityavkk wants to merge 1 commit into
electric-sql:mainfrom
adityavkk:fix/agents-server-stream-append-schema-lookup
Open

adityavkk wants to merge 1 commit into
electric-sql:mainfrom
adityavkk:fix/agents-server-stream-append-schema-lookup

Conversation

@adityavkk

Copy link
Copy Markdown

Fixes #4820

Problem

The stream-append route validates a typed append one event at a time (src/routing/stream-append.ts:120-135), and EntityManager.validateWriteEvent re-reads the entity type from Postgres on every call (getEffectiveSchemas -> registry.getEntityType, src/entity-manager.ts:3716 and :4035). An append of N events to a typed entity therefore costs 1 + N sequential registry round trips before it is forwarded, all N type reads returning the same row. With a local Postgres this is invisible; with Postgres ~60 ms away, a batch of ~10 events spends about 660 ms in the append path instead of about 120 ms. It also means the events of one append can be validated against different type revisions if the type is re-registered mid-batch.

Fix

Resolve the effective schemas once per append and validate every event in memory against that snapshot.

  • EntityManager.validateWriteEvents(entity, events) does the single getEffectiveSchemas lookup and returns the first invalid event's error, in event order, or null.
  • The per-event checks move verbatim into a private synchronous helper, validateEventAgainstStateSchemas(stateSchemas, event): unknown type -> 422 UNKNOWN_EVENT_TYPE; delete with old_value validates old_value, otherwise value; an undefined payload passes; a validator failure -> 422 SCHEMA_VALIDATION_FAILED.
  • validateWriteEvent(entity, event) becomes return this.validateWriteEvents(entity, [event]), so its signature and behaviour are unchanged for writeCollection and the existing test.
  • The route replaces the loop with one validateWriteEvents call. The write-token, fork-lock and stopped checks still run before validation; forwarding and the fire-and-forget side effects are untouched.

Round trips per typed append go from 1 + N to 2. Untyped entities and appends with no events still cost one round trip. No SQL, drizzle schema, registry API, or error code changes.

Trade-offs

  • It touches no SQL and keeps getEntityByStream and getEffectiveSchemas as they are (both have other callers), so the diff is small and the behaviour is easy to compare line by line.
  • No cache is introduced; every append still reads fresh entity and type state. The batch is now checked against one consistent schema snapshot rather than up to N reads of it.
  • A entities LEFT JOIN entity_types registry query would reach a single round trip, but it needs a new registry method duplicating getEntityByStream's /main parsing and getEffectiveSchemas' merge. Left as a possible follow-up.
  • writeCollection (src/entity-manager.ts:2466-2467, :2530) still reads the type twice (getEffectiveSchemas then validateWriteEvent). Deliberately not changed here to keep the diff minimal; it could call the new helper with the schemas it already holds.

Tests

New test/stream-append-route.test.ts drives the real electricAgentsStreamAppendRouter with a real EntityManager over a mocked registry (no Docker):

  • resolves the entity type once for a multi-event append: a 10-event typed append is forwarded once and registry.getEntityType is called once. Fails on main (called 10 times), passes with this change.
  • rejects the first invalid event in order without forwarding: event 2 has an unknown type and event 3 fails its schema; the response is 422 UNKNOWN_EVENT_TYPE for event 2 and forward is not called.
  • reports a schema failure at its position in the batch: event 3 fails its schema -> 422 SCHEMA_VALIDATION_FAILED, not forwarded.
  • skips validation for an untyped entity: forwarded, getEntityType never called.
  • validates every event of an append against one schema snapshot: the registry returns a different type revision on the second read; both observed_item events are accepted. Fails on main (the second event is checked against the second revision and rejected), passes with this change.

test/electric-agents-manager-write-validation.test.ts gains ElectricAgentsManager.validateWriteEvents / validates a batch against the entity's own schemas when the type row is missing (type lookup returns null, the entity's own state_schemas are used, one lookup). The existing validateWriteEvent delete/old_value test is unchanged and still passes.

Verification

pnpm install --frozen-lockfile
pnpm -r --filter "@electric-ax/agents-server^..." build

Red: on origin/main (140a5f4e7) with only the test files added:

$ cd packages/agents-server && pnpm exec vitest run test/stream-append-route.test.ts test/electric-agents-manager-write-validation.test.ts

 FAIL  test/stream-append-route.test.ts > stream append route > resolves the entity type once for a multi-event append
AssertionError: expected "vi.fn()" to be called 1 times, but got 10 times
 FAIL  test/stream-append-route.test.ts > stream append route > validates every event of an append against one schema snapshot
AssertionError: expected 422 to be 200 // Object.is equality
 FAIL  test/electric-agents-manager-write-validation.test.ts > ElectricAgentsManager.validateWriteEvents > validates a batch against the entity's own schemas when the type row is missing
TypeError: manager.validateWriteEvents is not a function

 Test Files  2 failed (2)
      Tests  3 failed | 18 passed (21)

Green: on this branch:

$ cd packages/agents-server && pnpm exec vitest run test/stream-append-route.test.ts test/electric-agents-manager-write-validation.test.ts

 ✓ test/stream-append-route.test.ts (5 tests) 32ms
 ✓ test/electric-agents-manager-write-validation.test.ts (16 tests) 31ms

 Test Files  2 passed (2)
      Tests  21 passed (21)
$ pnpm --filter @electric-ax/agents-server typecheck    # tsc --noEmit: clean
$ pnpm --filter @electric-ax/agents-server stylecheck   # eslint . --quiet: clean
$ pnpm exec prettier --check packages/agents-server/src/entity-manager.ts \
    packages/agents-server/src/routing/stream-append.ts \
    packages/agents-server/test/stream-append-route.test.ts \
    packages/agents-server/test/electric-agents-manager-write-validation.test.ts \
    .changeset/agents-server-append-schema-lookup.md
All matched files use Prettier code style!
$ git diff --check                                      # clean
$ GITHUB_BASE_REF=main node scripts/check-changeset.mjs
✅ Changesets cover all affected packages: @electric-ax/agents-server

The three red cases are the two route tests that encode the fix (lookup count, single snapshot) and the new manager test, whose method does not exist on main; the other 18 cases pass on both sides, which is the "semantics unchanged" evidence. The full agents-server suite needs the docker-compose Postgres + Electric backend and was not run here.

Files changed

  • packages/agents-server/src/entity-manager.ts: add validateWriteEvents and the private validateEventAgainstStateSchemas helper; validateWriteEvent delegates to the batch method.
  • packages/agents-server/src/routing/stream-append.ts: validate the whole append with one validateWriteEvents call instead of a per-event loop.
  • packages/agents-server/test/stream-append-route.test.ts: new route-level tests for the lookup count, error order, untyped entities and the single schema snapshot.
  • packages/agents-server/test/electric-agents-manager-write-validation.test.ts: unit test for the missing-type-row fallback of validateWriteEvents.
  • .changeset/agents-server-append-schema-lookup.md: patch changeset for @electric-ax/agents-server.

A typed stream append validated each event with validateWriteEvent,
which re-read the entity type from Postgres for every event, so an
append of N events cost 1 + N sequential registry round trips before
it was forwarded to the durable-streams server.

Resolve the effective schemas once per append and validate every event
in memory against that snapshot: the same checks and error codes in the
same order, the first invalid event reported, for two round trips
instead of 1 + N. validateWriteEvent keeps its single-event signature
for writeCollection.

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

agents-server: stream append reads the entity type from Postgres once per event (1 + N round trips per typed append)

1 participant