Skip to content

Commit 14a4b82

Browse files
authored
optimistic seperation (#260)
1 parent ebe32b9 commit 14a4b82

4 files changed

Lines changed: 154 additions & 109 deletions

File tree

.changeset/smart-parts-bake.md

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,5 @@
1+
---
2+
"@effect-rx/rx": patch
3+
---
4+
5+
seperate optimistic state from mutations

docs/rx/Rx.ts.md

Lines changed: 39 additions & 28 deletions
Original file line numberDiff line numberDiff line change
@@ -27,6 +27,9 @@ Added in v1.0.0
2727
- [windowFocusSignal](#windowfocussignal)
2828
- [KeyValueStore](#keyvaluestore)
2929
- [kvs](#kvs)
30+
- [Optimistic](#optimistic)
31+
- [optimistic](#optimistic-1)
32+
- [optimisticFn](#optimisticfn)
3033
- [Serializable](#serializable)
3134
- [Serializable (interface)](#serializable-interface)
3235
- [SerializableTypeId](#serializabletypeid)
@@ -48,7 +51,6 @@ Added in v1.0.0
4851
- [transform](#transform)
4952
- [withFallback](#withfallback)
5053
- [withLabel](#withlabel)
51-
- [withOptimisticSet](#withoptimisticset)
5254
- [constructors](#constructors)
5355
- [context](#context)
5456
- [family](#family)
@@ -234,6 +236,42 @@ export declare const kvs: <A>(options: {
234236
235237
Added in v1.0.0
236238
239+
# Optimistic
240+
241+
## optimistic
242+
243+
**Signature**
244+
245+
```ts
246+
export declare const optimistic: <A>(
247+
self: Rx<A>
248+
) => Writable<A, Rx<Result.Result<A extends Result.Result<infer _A, infer _E> ? _A : A, unknown>>>
249+
```
250+
251+
Added in v1.0.0
252+
253+
## optimisticFn
254+
255+
**Signature**
256+
257+
```ts
258+
export declare const optimisticFn: {
259+
<A, W, XA, XE, OW = A extends Result.Result<infer _A, infer _E> ? _A : A>(options: {
260+
readonly updateToValue: (value: OW, current: NoInfer<A>) => NoInfer<W>
261+
readonly fn: RxResultFn<NoInfer<OW>, XA, XE>
262+
}): (self: Writable<A, Rx<Result.Result<W, unknown>>>) => RxResultFn<OW, XA, XE>
263+
<A, W, XA, XE, OW = A extends Result.Result<infer _A, infer _E> ? _A : A>(
264+
self: Writable<A, Rx<Result.Result<W, unknown>>>,
265+
options: {
266+
readonly updateToValue: (value: OW, current: NoInfer<A>) => NoInfer<W>
267+
readonly fn: RxResultFn<NoInfer<OW>, XA, XE>
268+
}
269+
): RxResultFn<OW, XA, XE>
270+
}
271+
```
272+
273+
Added in v1.0.0
274+
237275
# Serializable
238276
239277
## Serializable (interface)
@@ -493,33 +531,6 @@ export declare const withLabel: {
493531
494532
Added in v1.0.0
495533
496-
## withOptimisticSet
497-
498-
**Signature**
499-
500-
```ts
501-
export declare const withOptimisticSet: {
502-
<A, XA, XE, W = A extends Result.Result<infer _A, infer _E> ? _A : A>(options: {
503-
readonly updateToValue: (
504-
value: W,
505-
current: NoInfer<A>
506-
) => A extends Result.Result<infer _A, infer _E> ? _A : NoInfer<A>
507-
readonly fn: RxResultFn<NoInfer<W>, XA, XE>
508-
readonly disableRefresh?: boolean | undefined
509-
}): (self: Rx<A>) => Writable<A, W>
510-
<A, XA, XE, W = A extends Result.Result<infer _A, infer _E> ? _A : A>(
511-
self: Rx<A>,
512-
options: {
513-
readonly updateToValue: (value: W, current: NoInfer<A>) => A extends Result.Result<infer _A, infer _E> ? _A : A
514-
readonly fn: RxResultFn<NoInfer<W>, XA, XE>
515-
readonly disableRefresh?: boolean | undefined
516-
}
517-
): Writable<A, W>
518-
}
519-
```
520-
521-
Added in v1.0.0
522-
523534
# constructors
524535
525536
## context

packages/rx/src/Rx.ts

Lines changed: 95 additions & 73 deletions
Original file line numberDiff line numberDiff line change
@@ -1355,88 +1355,110 @@ export const debounce: {
13551355

13561356
/**
13571357
* @since 1.0.0
1358-
* @category combinators
1358+
* @category Optimistic
13591359
*/
1360-
export const withOptimisticSet: {
1361-
<A, XA, XE, W = A extends Result.Result<infer _A, infer _E> ? _A : A>(
1360+
export const optimistic = <A>(
1361+
self: Rx<A>
1362+
): Writable<A, Rx<Result.Result<A extends Result.Result<infer _A, infer _E> ? _A : A, unknown>>> => {
1363+
let counter = 0
1364+
const writeRx = state(
1365+
[
1366+
counter,
1367+
undefined as any as Rx<Result.Result<A extends Result.Result<infer _A, infer _E> ? _A : A, unknown>>
1368+
] as const
1369+
)
1370+
return writable(
1371+
(get) => {
1372+
let lastValue = get.once(self)
1373+
const isResult = Result.isResult(lastValue)
1374+
get.subscribe(self, (value) => {
1375+
lastValue = value
1376+
if (!Result.isResult(value)) {
1377+
return get.setSelf(value)
1378+
}
1379+
const current = Option.getOrUndefined(get.self<Result.Result<any, any>>())!
1380+
if (Result.isSuccess(current) && Result.isSuccess(value)) {
1381+
if (value.timestamp >= current.timestamp) {
1382+
get.setSelf(value)
1383+
}
1384+
} else {
1385+
get.setSelf(value)
1386+
}
1387+
})
1388+
const transitions = new Set<Rx<Result.Result<A extends Result.Result<infer _A, infer _E> ? _A : A, unknown>>>()
1389+
const cancels = new Set<() => void>()
1390+
get.subscribe(writeRx, ([, rx]) => {
1391+
if (transitions.has(rx)) {
1392+
return
1393+
}
1394+
transitions.add(rx)
1395+
const cancel = get.registry.subscribe(rx, (result) => {
1396+
if (Result.isSuccess(result) && result.waiting) {
1397+
return get.setSelf(isResult ? Result.success(result.value, { waiting: true }) : result.value)
1398+
}
1399+
transitions.delete(rx)
1400+
cancels.delete(cancel)
1401+
cancel()
1402+
if (transitions.size === 0) {
1403+
if (Result.isFailure(result)) {
1404+
get.setSelf(lastValue)
1405+
}
1406+
get.refresh(self)
1407+
}
1408+
}, { immediate: true })
1409+
cancels.add(cancel)
1410+
})
1411+
get.addFinalizer(() => {
1412+
for (const cancel of cancels) {
1413+
cancel()
1414+
}
1415+
transitions.clear()
1416+
cancels.clear()
1417+
})
1418+
return lastValue
1419+
},
1420+
(ctx, rx) => ctx.set(writeRx, [++counter, rx]),
1421+
(refresh) => refresh(self)
1422+
)
1423+
}
1424+
1425+
/**
1426+
* @since 1.0.0
1427+
* @category Optimistic
1428+
*/
1429+
export const optimisticFn: {
1430+
<A, W, XA, XE, OW = A extends Result.Result<infer _A, infer _E> ? _A : A>(
13621431
options: {
1363-
readonly updateToValue: (
1364-
value: W,
1365-
current: NoInfer<A>
1366-
) => A extends Result.Result<infer _A, infer _E> ? _A : NoInfer<A>
1367-
readonly fn: RxResultFn<NoInfer<W>, XA, XE>
1368-
readonly disableRefresh?: boolean | undefined
1432+
readonly updateToValue: (value: OW, current: NoInfer<A>) => NoInfer<W>
1433+
readonly fn: RxResultFn<NoInfer<OW>, XA, XE>
13691434
}
13701435
): (
1371-
self: Rx<A>
1372-
) => Writable<A, W>
1373-
<A, XA, XE, W = A extends Result.Result<infer _A, infer _E> ? _A : A>(
1374-
self: Rx<A>,
1436+
self: Writable<A, Rx<Result.Result<W, unknown>>>
1437+
) => RxResultFn<OW, XA, XE>
1438+
<A, W, XA, XE, OW = A extends Result.Result<infer _A, infer _E> ? _A : A>(
1439+
self: Writable<A, Rx<Result.Result<W, unknown>>>,
13751440
options: {
1376-
readonly updateToValue: (
1377-
value: W,
1378-
current: NoInfer<A>
1379-
) => A extends Result.Result<infer _A, infer _E> ? _A : A
1380-
readonly fn: RxResultFn<NoInfer<W>, XA, XE>
1381-
readonly disableRefresh?: boolean | undefined
1441+
readonly updateToValue: (value: OW, current: NoInfer<A>) => NoInfer<W>
1442+
readonly fn: RxResultFn<NoInfer<OW>, XA, XE>
13821443
}
1383-
): Writable<A, W>
1384-
} = dual(2, <A, W, XA, XE>(
1385-
self: Rx<A>,
1444+
): RxResultFn<OW, XA, XE>
1445+
} = dual(2, <A, W, OW, XA, XE>(
1446+
self: Writable<A, Rx<Result.Result<W, unknown>>>,
13861447
options: {
1387-
readonly updateToValue: (value: W, current: A) => A extends Result.Result<infer _A, infer _E> ? _A : A
1388-
readonly fn: RxResultFn<W, XA, XE>
1389-
readonly disableRefresh?: boolean | undefined
1448+
readonly updateToValue: (value: OW, current: A) => W
1449+
readonly fn: RxResultFn<OW, XA, XE>
13901450
}
1391-
): Writable<A, W> => {
1392-
let counter = 0
1393-
const argRx = state([counter, undefined as A | undefined] as const)
1394-
return writable((get) => {
1395-
let lastValue = get.once(self)
1396-
get.subscribe(self, (value) => {
1397-
lastValue = value
1398-
if (!Result.isResult(value)) {
1399-
return get.setSelf(value)
1400-
}
1401-
const current = Option.getOrUndefined(get.self<Result.Result<any, any>>())!
1402-
if (Result.isSuccess(current) && Result.isSuccess(value)) {
1403-
if (value.timestamp >= current.timestamp) {
1404-
get.setSelf(value)
1405-
}
1406-
} else {
1407-
get.setSelf(lastValue)
1408-
}
1409-
})
1410-
let lastSetSuccess: A | undefined
1411-
get.subscribe(argRx, ([, arg]) => {
1412-
if (arg === undefined) {
1413-
return
1414-
}
1415-
lastSetSuccess = arg
1416-
get.setSelf(arg)
1417-
})
1418-
get.subscribe(options.fn, (value) => {
1419-
if (value.waiting || Result.isInitial(value)) {
1420-
return
1421-
} else if (Result.isFailure(value)) {
1422-
return get.setSelf(lastValue)
1423-
}
1424-
if (options.disableRefresh !== true) {
1425-
return get.refresh(self)
1426-
}
1427-
const current = Option.getOrUndefined(get.self<Result.Success<any, any>>())!
1428-
if (current === lastSetSuccess) {
1429-
get.setSelf(Result.success(current.value))
1430-
}
1451+
): RxResultFn<OW, XA, XE> => {
1452+
const transition = state<Result.Result<W, unknown>>(Result.initial())
1453+
return fn((arg: OW, get) => {
1454+
const value = options.updateToValue(arg, get(self))
1455+
get.set(transition, Result.success(value, { waiting: true }))
1456+
get.set(self, transition)
1457+
get.set(options.fn, arg)
1458+
return Effect.onExit(get.result(options.fn, { suspendOnWaiting: true }), (exit) => {
1459+
get.set(transition, Result.fromExit(Exit.as(exit, value)))
1460+
return Effect.void
14311461
})
1432-
return lastValue
1433-
}, (ctx, value) => {
1434-
const current = ctx.get(self)
1435-
const arg = options.updateToValue(value, current)
1436-
ctx.set(argRx, [++counter, Result.isResult(current) ? Result.success(arg, { waiting: true }) as A : arg as A])
1437-
ctx.set(options.fn, value)
1438-
}, (refresh) => {
1439-
refresh(self)
14401462
})
14411463
})
14421464

packages/rx/test/Rx.test.ts

Lines changed: 15 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -1045,14 +1045,15 @@ describe("Rx", () => {
10451045
})
10461046
})
10471047

1048-
describe("withOptimisticSet", () => {
1048+
describe("optimistic", () => {
10491049
it("non-Result", async () => {
10501050
const latch = Effect.unsafeMakeLatch()
10511051
const r = Registry.make()
10521052
let i = 0
10531053
const rx = Rx.make(() => i)
1054-
const optimisticRx = rx.pipe(
1055-
Rx.withOptimisticSet({
1054+
const optimisticRx = rx.pipe(Rx.optimistic)
1055+
const fn = optimisticRx.pipe(
1056+
Rx.optimisticFn({
10561057
updateToValue: (value) => value,
10571058
fn: Rx.fn(Effect.fnUntraced(function*() {
10581059
yield* latch.await
@@ -1063,7 +1064,7 @@ describe("Rx", () => {
10631064

10641065
expect(r.get(rx)).toEqual(0)
10651066
expect(r.get(optimisticRx)).toEqual(0)
1066-
r.set(optimisticRx, 1)
1067+
r.set(fn, 1)
10671068
i = 2
10681069

10691070
// optimistic phase: the optimistic value is set, but the true value is not
@@ -1084,7 +1085,10 @@ describe("Rx", () => {
10841085
let i = 0
10851086
const rx = Rx.make(Effect.sync(() => i))
10861087
const optimisticRx = rx.pipe(
1087-
Rx.withOptimisticSet({
1088+
Rx.optimistic
1089+
)
1090+
const fn = optimisticRx.pipe(
1091+
Rx.optimisticFn({
10881092
updateToValue: (value) => value,
10891093
fn: Rx.fn(Effect.fnUntraced(function*() {
10901094
yield* latch.await
@@ -1095,7 +1099,7 @@ describe("Rx", () => {
10951099

10961100
expect(r.get(rx)).toEqual(Result.success(0))
10971101
expect(r.get(optimisticRx)).toEqual(Result.success(0))
1098-
r.set(optimisticRx, 1)
1102+
r.set(fn, 1)
10991103
i = 2
11001104

11011105
// optimistic phase: the optimistic value is set, but the true value is not
@@ -1116,7 +1120,10 @@ describe("Rx", () => {
11161120
const i = 0
11171121
const rx = Rx.make(() => i)
11181122
const optimisticRx = rx.pipe(
1119-
Rx.withOptimisticSet({
1123+
Rx.optimistic
1124+
)
1125+
const fn = optimisticRx.pipe(
1126+
Rx.optimisticFn({
11201127
updateToValue: (value) => value,
11211128
fn: Rx.fn(Effect.fnUntraced(function*() {
11221129
yield* latch.await
@@ -1128,7 +1135,7 @@ describe("Rx", () => {
11281135

11291136
expect(r.get(rx)).toEqual(0)
11301137
expect(r.get(optimisticRx)).toEqual(0)
1131-
r.set(optimisticRx, 1)
1138+
r.set(fn, 1)
11321139

11331140
// optimistic phase: the optimistic value is set, but the true value is not
11341141
expect(r.get(rx)).toEqual(0)

0 commit comments

Comments
 (0)