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
6 changes: 6 additions & 0 deletions .changeset/tidy-derived-dependencies.md
Original file line number Diff line number Diff line change
@@ -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.
18 changes: 13 additions & 5 deletions packages/effect-app/src/DataDependencies.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<typeof DataDependency>
Expand Down Expand Up @@ -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))
}
Expand All @@ -57,7 +57,6 @@ const containsDependency = (dependencies: ReadonlySet<DataDependency>, 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))
Expand All @@ -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>): 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*() {
Expand Down Expand Up @@ -125,7 +132,8 @@ export const makeDataDependencyRecorder = (

export const repo = (name: string, ids?: NonEmptyReadonlyArray<string>): 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<string>): DataDependency =>
ids === undefined ? { type: "signal", name } : { type: "signal", name, ids }

const RepoReadScope = Context.Reference<DataDependency | undefined>("effect-app/DataDependencies/RepoReadScope", {
defaultValue: () => undefined
Expand Down
50 changes: 41 additions & 9 deletions packages/effect-app/src/Model/Repository/internal/internal.ts
Original file line number Diff line number Diff line change
Expand Up @@ -115,6 +115,7 @@ export function makeRepoInternal<
schemaContext?: Context.Context<RCtx>
makeInitial?: Effect.Effect<readonly T[], E, RInitial> | undefined
dependencyIds?: (item: T) => NonEmptyReadonlyArray<string>
additionalWriteDependencies?: (item: T) => ReadonlyArray<DataDependencies.DataDependency>
config?: Omit<StoreConfig<Encoded>, "partitionValue"> & {
partitionValue?: (e?: Encoded) => string
}
Expand All @@ -124,6 +125,7 @@ export function makeRepoInternal<
publishEvents: (evt: NonEmptyReadonlyArray<Evt>) => Effect.Effect<void, never, RPublish>
makeInitial?: Effect.Effect<readonly T[], E, RInitial> | undefined
dependencyIds?: (item: T) => NonEmptyReadonlyArray<string>
additionalWriteDependencies?: (item: T) => ReadonlyArray<DataDependencies.DataDependency>
config?: Omit<StoreConfig<Encoded>, "partitionValue"> & {
partitionValue?: (e?: Encoded) => string
}
Expand All @@ -148,18 +150,28 @@ export function makeRepoInternal<
const recordRead = DataDependencies.readRepo(name)
const entityDependency = (ids: NonEmptyReadonlyArray<T[IdKey]>) =>
DataDependencies.repo(name, [String(ids[0]), ...ids.slice(1).map(String)])
const itemDependency = (items: NonEmptyReadonlyArray<T>) => {
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<T[IdKey]>) =>
DataDependencies.write(entityDependency(ids))
const recordItemWrite = (items: NonEmptyReadonlyArray<T>) => DataDependencies.write(itemDependency(items))
const recordItemWrite = (items: NonEmptyReadonlyArray<T>) =>
Effect.forEach(
DataDependencies.merge(new Set(items.flatMap(itemDependencies))),
DataDependencies.write,
{ discard: true }
)
const recordAdditionalItemWrites = (items: ReadonlyArray<T>) =>
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)
Expand Down Expand Up @@ -292,6 +304,20 @@ export function makeRepoInternal<
)
})

const loadExistingItems = (ids: ReadonlyArray<T[IdKey]>) =>
Effect
.forEach(
ids,
(id) =>
findE(id).pipe(
Effect.flatMap(Option.match({
onNone: () => Effect.succeed(Option.none<T>()),
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 }
Expand Down Expand Up @@ -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({
Expand Down Expand Up @@ -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 })
Expand Down
8 changes: 8 additions & 0 deletions packages/effect-app/src/Model/Repository/makeRepo.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down Expand Up @@ -61,6 +62,13 @@ export interface RepositoryOptions<
/** IDs, including aliases, that identify an entity for dependency invalidation. */
dependencyIds?: (item: T) => NonEmptyReadonlyArray<string>

/**
* 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<DataDependencies.DataDependency>

overrides?: (
repo: Repository<T, Encoded, Evt, ItemType, IdKey, Exclude<RSchema, RCtx>, RPublish, RCtx>
) => Repository<T, Encoded, Evt, ItemType, IdKey, Exclude<RSchema, RCtx>, RPublish, RCtx>
Expand Down
75 changes: 75 additions & 0 deletions packages/infra/test/repository-ext.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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*() {
Expand Down
Loading