-
Notifications
You must be signed in to change notification settings - Fork 70
Expand file tree
/
Copy pathSequence.fs
More file actions
75 lines (57 loc) · 2.82 KB
/
Copy pathSequence.fs
File metadata and controls
75 lines (57 loc) · 2.82 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
// Manages a sequence of ids, without provision for returning unused ones in cases where we're potentially leaving a gap
// see Gapless.fs for a potential approach for handling such a desire
module Sequence
open System
let [<Literal>] Category = "Sequence"
let streamName id = FsCodec.StreamName.create Category (SequenceId.toString id)
// shim for net461
module Seq =
let tryLast (source : seq<_>) =
use e = source.GetEnumerator()
if e.MoveNext() then
let mutable res = e.Current
while (e.MoveNext()) do res <- e.Current
Some res
else
None
// NOTE - these types and the union case names reflect the actual storage formats and hence need to be versioned with care
module Events =
type Reserved = { next : int64 }
type Event =
| Reserved of Reserved
interface TypeShape.UnionContract.IUnionContract
let codecNewtonsoft = FsCodec.NewtonsoftJson.Codec.Create<Event>()
let codecStj = FsCodec.SystemTextJson.Codec.Create<Event>()
module Fold =
type State = { next : int64 }
let initial = { next = 0L }
let private evolve _ignoreState = function
| Events.Reserved e -> { next = e.next }
let fold (state: State) (events: seq<Events.Event>) : State =
Seq.tryLast events |> Option.fold evolve state
let snapshot (state : State) = Events.Reserved { next = state.next }
let decideReserve (count : int) (state : Fold.State) : int64 * Events.Event list =
state.next,[Events.Reserved { next = state.next + int64 count }]
type Service internal (resolve : SequenceId -> Equinox.Stream<Events.Event, Fold.State>) =
/// Reserves an id, yielding the reserved value. Optional <c>count</c> enables reserving more than the default count of <c>1</c> in a single transaction
member __.Reserve(series,?count) : Async<int64> =
let stream = resolve series
stream.Transact(decideReserve (defaultArg count 1))
let create resolve =
let resolve sequenceId =
let streamName = streamName sequenceId
Equinox.Stream(Serilog.Log.ForContext<Service>(), resolve streamName, maxAttempts = 3)
Service(resolve)
module Cosmos =
open Equinox.CosmosStore
let private create (context,cache,accessStrategy) =
let cacheStrategy = CachingStrategy.SlidingWindow (cache, TimeSpan.FromMinutes 20.) // OR CachingStrategy.NoCaching
let category = CosmosStoreCategory(context, Events.codecStj, Fold.fold, Fold.initial, cacheStrategy, accessStrategy)
create category.Resolve
module LatestKnownEvent =
let create (context,cache) =
let accessStrategy = AccessStrategy.LatestKnownEvent
create (context,cache,accessStrategy)
module RollingUnfolds =
let create (context,cache) =
create (context,cache,AccessStrategy.RollingState Fold.snapshot)