From 8cf52c815e30f73067d0d2e8c7163a7d40c4f430 Mon Sep 17 00:00:00 2001 From: "omegent-app[bot]" <306514130+omegent-app[bot]@users.noreply.github.com> Date: Mon, 3 Aug 2026 21:49:09 +0000 Subject: [PATCH] feat(repository): derive affected query dependencies --- .changeset/tidy-derived-dependencies.md | 6 ++ packages/effect-app/src/DataDependencies.ts | 18 +++-- .../src/Model/Repository/internal/internal.ts | 50 ++++++++++--- .../src/Model/Repository/makeRepo.ts | 8 ++ packages/infra/test/repository-ext.test.ts | 75 +++++++++++++++++++ 5 files changed, 143 insertions(+), 14 deletions(-) create mode 100644 .changeset/tidy-derived-dependencies.md diff --git a/.changeset/tidy-derived-dependencies.md b/.changeset/tidy-derived-dependencies.md new file mode 100644 index 000000000..617b9570d --- /dev/null +++ b/.changeset/tidy-derived-dependencies.md @@ -0,0 +1,6 @@ +--- +"effect-app": minor +"@effect-app/infra": patch +--- + +Support ID-scoped signal dependencies and additive repository dependencies derived from previous and current entities. diff --git a/packages/effect-app/src/DataDependencies.ts b/packages/effect-app/src/DataDependencies.ts index 8b6b8fbdc..920fc7cb2 100644 --- a/packages/effect-app/src/DataDependencies.ts +++ b/packages/effect-app/src/DataDependencies.ts @@ -14,7 +14,8 @@ export const DataDependency = S.Union([ }), S.Struct({ type: S.Literal("signal"), - name: S.String + name: S.String, + ids: S.optional(S.NonEmptyArray(S.String)) }) ]) export type DataDependency = S.Schema.Type @@ -42,7 +43,6 @@ const sameDomain = (a: DataDependency, b: DataDependency) => a.type === b.type & const dependenciesOverlap = (a: DataDependency, b: DataDependency) => { if (!sameDomain(a, b)) return false - if (a.type === "signal" || b.type === "signal") return true if (a.ids === undefined || b.ids === undefined) return true return a.ids.some((id) => b.ids?.includes(id)) } @@ -57,7 +57,6 @@ const containsDependency = (dependencies: ReadonlySet, dependenc const appendDependency = (dependency: DataDependency) => (dependencies: DataDependencies): DataDependencies => { const existing = [...dependencies].find((_) => sameDomain(_, dependency)) if (existing === undefined) return new Set([...dependencies, dependency]) - if (existing.type === "signal" || dependency.type === "signal") return dependencies if (existing.ids === undefined) return dependencies if (dependency.ids === undefined) { return new Set([...dependencies].filter((_) => _ !== existing).concat(dependency)) @@ -68,10 +67,18 @@ const appendDependency = (dependency: DataDependency) => (dependencies: DataDepe ...new Set([...existing.ids.slice(1), ...dependency.ids].filter((id) => id !== first)) ] if (ids.length === existing.ids.length) return dependencies - const merged = repo(existing.name, ids) + const merged = existing.type === "repo" ? repo(existing.name, ids) : signal(existing.name, ids) return new Set([...dependencies].filter((_) => _ !== existing).concat(merged)) } +export const merge = (...groups: ReadonlyArray): DataDependencies => { + let merged = empty() + for (const group of groups) { + for (const dependency of group) merged = appendDependency(dependency)(merged) + } + return merged +} + export const DataDependencyRecorder = RequestScopedDependencies.make( "effect-app/DataDependencyRecorder", Effect.gen(function*() { @@ -125,7 +132,8 @@ export const makeDataDependencyRecorder = ( export const repo = (name: string, ids?: NonEmptyReadonlyArray): DataDependency => ids === undefined ? { type: "repo", name } : { type: "repo", name, ids } -export const signal = (name: string): DataDependency => ({ type: "signal", name }) +export const signal = (name: string, ids?: NonEmptyReadonlyArray): DataDependency => + ids === undefined ? { type: "signal", name } : { type: "signal", name, ids } const RepoReadScope = Context.Reference("effect-app/DataDependencies/RepoReadScope", { defaultValue: () => undefined diff --git a/packages/effect-app/src/Model/Repository/internal/internal.ts b/packages/effect-app/src/Model/Repository/internal/internal.ts index 835d14cd6..b0c3d4c87 100644 --- a/packages/effect-app/src/Model/Repository/internal/internal.ts +++ b/packages/effect-app/src/Model/Repository/internal/internal.ts @@ -115,6 +115,7 @@ export function makeRepoInternal< schemaContext?: Context.Context makeInitial?: Effect.Effect | undefined dependencyIds?: (item: T) => NonEmptyReadonlyArray + additionalWriteDependencies?: (item: T) => ReadonlyArray config?: Omit, "partitionValue"> & { partitionValue?: (e?: Encoded) => string } @@ -124,6 +125,7 @@ export function makeRepoInternal< publishEvents: (evt: NonEmptyReadonlyArray) => Effect.Effect makeInitial?: Effect.Effect | undefined dependencyIds?: (item: T) => NonEmptyReadonlyArray + additionalWriteDependencies?: (item: T) => ReadonlyArray config?: Omit, "partitionValue"> & { partitionValue?: (e?: Encoded) => string } @@ -148,18 +150,28 @@ export function makeRepoInternal< const recordRead = DataDependencies.readRepo(name) const entityDependency = (ids: NonEmptyReadonlyArray) => DataDependencies.repo(name, [String(ids[0]), ...ids.slice(1).map(String)]) - const itemDependency = (items: NonEmptyReadonlyArray) => { - const firstIds = args.dependencyIds?.(items[0]) ?? [String(items[0][idKey])] - return DataDependencies.repo(name, [ - firstIds[0], - ...firstIds.slice(1), - ...items.slice(1).flatMap((item) => args.dependencyIds?.(item) ?? [String(item[idKey])]) - ]) + const itemDependencies = (item: T) => { + const ids = args.dependencyIds?.(item) ?? [String(item[idKey])] + return [ + DataDependencies.repo(name, ids), + ...(args.additionalWriteDependencies?.(item) ?? []) + ] } const recordEntityRead = (id: T[IdKey]) => DataDependencies.read(entityDependency([id])) const recordEntityWrite = (ids: NonEmptyReadonlyArray) => DataDependencies.write(entityDependency(ids)) - const recordItemWrite = (items: NonEmptyReadonlyArray) => DataDependencies.write(itemDependency(items)) + const recordItemWrite = (items: NonEmptyReadonlyArray) => + Effect.forEach( + DataDependencies.merge(new Set(items.flatMap(itemDependencies))), + DataDependencies.write, + { discard: true } + ) + const recordAdditionalItemWrites = (items: ReadonlyArray) => + Effect.forEach( + DataDependencies.merge(new Set(items.flatMap((item) => args.additionalWriteDependencies?.(item) ?? []))), + DataDependencies.write, + { discard: true } + ) const cms = Effect.map(getContextMap.pipe(Effect.orDie), (_) => ({ get: (id: string) => _.get(`${name}.${id}`), set: (id: string, etag: string | undefined) => _.set(`${name}.${id}`, etag) @@ -292,6 +304,20 @@ export function makeRepoInternal< ) }) + const loadExistingItems = (ids: ReadonlyArray) => + Effect + .forEach( + ids, + (id) => + findE(id).pipe( + Effect.flatMap(Option.match({ + onNone: () => Effect.succeed(Option.none()), + onSome: (encoded) => decode(encoded).pipe(Effect.orDie, Effect.map(Option.some)) + })) + ) + ) + .pipe(Effect.map(Array.getSomes)) + const find = Effect.fn("Repository.find", { kind: "client", attributes: { "app.entity": name } @@ -327,7 +353,10 @@ export function makeRepoInternal< const it = Chunk.fromIterable(items) if (Chunk.isNonEmpty(it)) { const values = Chunk.toReadonlyArray(it) - yield* recordItemWrite([values[0], ...values.slice(1)]) + const previous = args.additionalWriteDependencies + ? yield* loadExistingItems(values.map((item) => item[idKey])) + : [] + yield* recordItemWrite([values[0], ...values.slice(1), ...previous]) } const evts = [...events] yield* Effect.annotateCurrentSpan({ @@ -385,6 +414,9 @@ export function makeRepoInternal< return } yield* recordEntityWrite(ids) + if (args.additionalWriteDependencies) { + yield* loadExistingItems(ids).pipe(Effect.flatMap(recordAdditionalItemWrites)) + } const { set } = yield* cms const eids = yield* Effect.forEach(ids, (_) => encodeIdOnly(_ as any)).pipe(Effect.orDie) yield* Effect.annotateCurrentSpan({ "app.entity.ids": eids }) diff --git a/packages/effect-app/src/Model/Repository/makeRepo.ts b/packages/effect-app/src/Model/Repository/makeRepo.ts index 8739b3810..b2b5af115 100644 --- a/packages/effect-app/src/Model/Repository/makeRepo.ts +++ b/packages/effect-app/src/Model/Repository/makeRepo.ts @@ -9,6 +9,7 @@ import type * as Scope from "effect/Scope" import type { NonEmptyReadonlyArray } from "../../Array.ts" import type * as Context from "../../Context.ts" +import type * as DataDependencies from "../../DataDependencies.ts" import * as Effect from "../../Effect.ts" import type * as S from "../../Schema.ts" import type { StoreConfig, StoreMaker } from "../../Store.ts" @@ -61,6 +62,13 @@ export interface RepositoryOptions< /** IDs, including aliases, that identify an entity for dependency invalidation. */ dependencyIds?: (item: T) => NonEmptyReadonlyArray + /** + * Additional derived queries or resources affected by writing an item. + * Saves invalidate dependencies derived from both the previous and next item; + * removals derive them from the removed item. + */ + additionalWriteDependencies?: (item: T) => ReadonlyArray + overrides?: ( repo: Repository, RPublish, RCtx> ) => Repository, RPublish, RCtx> diff --git a/packages/infra/test/repository-ext.test.ts b/packages/infra/test/repository-ext.test.ts index 07579a919..e5f1201af 100644 --- a/packages/infra/test/repository-ext.test.ts +++ b/packages/infra/test/repository-ext.test.ts @@ -245,6 +245,81 @@ describe("repository ext save/remove batching", () => { Effect.provide(TestStoreLive) )) + it.effect("records repository and affected-query write dependencies", () => + Effect + .gen(function*() { + const readsRef = yield* Ref.make(DataDependencies.empty()) + const writesRef = yield* Ref.make(DataDependencies.empty()) + const recorder = DataDependencies.makeDataDependencyRecorder(readsRef, writesRef) + + yield* Effect + .gen(function*() { + const repo = yield* makeRepo("DependencyItem", BatchItem, { + dependencyIds: (item) => [item.id, `alias-${item.id}`], + additionalWriteDependencies: (item) => [ + DataDependencies.signal("DependencyItem.List", [item.label]) + ] + }) + yield* repo.save(new BatchItem({ id: "1", label: "one" })) + yield* DataDependencies.read(DataDependencies.signal("DependencyItem.List", ["one"])) + }) + .pipe(Effect.provideService(DataDependencies.DataDependencyRecorder, recorder)) + + expect(yield* Ref.get(readsRef)).toEqual( + new Set([DataDependencies.signal("DependencyItem.List", ["one"])]) + ) + expect(yield* Ref.get(writesRef)).toEqual( + new Set([ + DataDependencies.repo("DependencyItem", ["1", "alias-1"]), + DataDependencies.signal("DependencyItem.List", ["one"]) + ]) + ) + }) + .pipe( + setupRequestContextFromCurrent(), + Effect.provide(TestStoreLive) + )) + + it.effect("invalidates derived dependencies from previous saves and removeById", () => + Effect + .gen(function*() { + const readsRef = yield* Ref.make(DataDependencies.empty()) + const writesRef = yield* Ref.make(DataDependencies.empty()) + const recorder = DataDependencies.makeDataDependencyRecorder(readsRef, writesRef) + + yield* Effect + .gen(function*() { + const repo = yield* makeRepo("DerivedDependencyItem", BatchItem, { + additionalWriteDependencies: (item) => [ + DataDependencies.signal("DerivedDependencyItem.List", [item.label]) + ] + }) + yield* repo.save(new BatchItem({ id: "1", label: "old" })) + yield* recorder.drainWrites + + yield* repo.save(new BatchItem({ id: "1", label: "new" })) + expect(yield* recorder.drainWrites).toEqual( + new Set([ + DataDependencies.repo("DerivedDependencyItem", ["1"]), + DataDependencies.signal("DerivedDependencyItem.List", ["new", "old"]) + ]) + ) + + yield* repo.removeById("1") + expect(yield* recorder.drainWrites).toEqual( + new Set([ + DataDependencies.repo("DerivedDependencyItem", ["1"]), + DataDependencies.signal("DerivedDependencyItem.List", ["new"]) + ]) + ) + }) + .pipe(Effect.provideService(DataDependencies.DataDependencyRecorder, recorder)) + }) + .pipe( + setupRequestContextFromCurrent(), + Effect.provide(TestStoreLive) + )) + it.effect("records schema timing on repository spans without codec child spans", () => Effect .gen(function*() {