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
3 changes: 2 additions & 1 deletion CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -17,7 +17,7 @@ The `Unreleased` section name is replaced by the expected version of next releas
- `Cosmos`: Reorganize Sync log message text, merge with Sync Conflict message [#241](https://github.com/jet/equinox/pull/241)
- `Cosmos`: Converge Stored Procedure Impl with `tip-isa-batch` impl from V3 (minor Request Charges cost reduction) [#242](https://github.com/jet/equinox/pull/242)

- Fork `Equinox.Cosmos` to `Equinox.CosmosStore`
- Fork `Equinox.Cosmos` to `Equinox.CosmosStore`:
- target `Microsoft.Azure.Cosmos` v `3.9.0` (instead of `Microsoft.Azure.DocumentDB`[`.Core`] v 2.x) [#144](https://github.com/jet/equinox/pull/144)
- Removed [warmup call](https://github.com/Azure/azure-cosmos-dotnet-v3/issues/1436)
- Rename `Equinox.Cosmos` DLL and namespace to `Equinox.CosmosStore` [#243](https://github.com/jet/equinox/pull/243)
Expand All @@ -31,6 +31,7 @@ The `Unreleased` section name is replaced by the expected version of next releas
- Rename `Equinox.Cosmos.Resolver` -> `Equinox.CosmosStore.CosmosStoreCategory`
- Rename `Equinox.Cosmos.Connector` -> `Equinox.CosmosStore.CosmosStoreClientFactory`
- Reorganized `QueryRetryPolicy` to handle `IAsyncEnumerable` coming in Cosmos SDK V4 [#246](https://github.com/jet/equinox/pull/246) :pray: [@ylibrach](https://github.com/ylibrach)
- Added Secondary store fallback for Event loading, enabling Streams to be hot-migrated (archived to a secondary/clone, then pruned from the primary/active) between Primary and Secondary stores [#247](https://github.com/jet/equinox/pull/247)
- target `EventStore.Client` v `20.6` (instead of v `5.0.x`) [#224](https://github.com/jet/equinox/pull/224)
- Retarget `netcoreapp2.1` apps to `netcoreapp3.1` with `SystemTextJson`
- Retarget Todobackend to `aspnetcore` v `3.1`
Expand Down
54 changes: 50 additions & 4 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -303,7 +303,6 @@ While Equinox is implemented in F#, and F# is a great fit for writing event-sour
# cosmos specifies source overrides (using defaults in step 1 in this instance)
dotnet run -- -g projector2 cosmos
```

7. Use `propulsion` tool to Run a CosmosDb ChangeFeedProcessor, emitting to a Kafka topic

```powershell
Expand Down Expand Up @@ -341,8 +340,51 @@ While Equinox is implemented in F#, and F# is a great fit for writing event-sour
dotnet run -- -t topic0 -g consumer1
```

9. Generate an Archive container; Generate a ChangeFeedProcessor App to mirror desired streams from the Primary to it

```powershell
# once
eqx init -ru 400 cosmos -c equinox-test-archive

md archiver | cd

# Generate a template app that'll sync from the Primary (i.e. equinox-test)
# to the Secondary (i.e. equinox-test-archive)
dotnet new proArchiver

# TODO edit Handler.fs to add criteria for what to Archive
# - Normally you won't want to Archive stuff like e.g. `Sync-` checkppoint streams
# - Any other ephemeral application streams can be excluded too

# -w 4 # constrain parallel writers in order to leave headroom for readers; Secondary container should be cheaper to run
# -S -t 40 # emit log messages for Sync calls costing > 40 RU
# -md 20 (or lower) is recommended to be nice to the writers - the archiver can afford to lag
dotnet run -c Release -- -w 4 -S -t 40 -g ArchiverConsumer `
cosmos -md 20 -c equinox-test -a equinox-test-aux `
cosmos -c equinox-test-archive
```

10. Use a ChangeFeedProcessor driven from the Archive Container to Prune the Primary

```powershell
md pruner | cd

# Generate a template app that'll read from the Archive (i.e. equinox-test-archive)
# and prune expired events from the Primary (i.e. equinox-test)
dotnet new proPruner

# TODO edit Handler.fs to add criteria for what to Prune
# - While its possible to prune the minute it's archived, normally you'll want to allow a time lag before doing so

# -w 2 # constrain parallel pruners in order to not consume RUs excessively on Primary
# -md 10 (or lower) is recommended to contrain consumption on the Secondary - Pruners lagging is rarely critical
dotnet run -c Release -- -w 2 -g PrunerConsumer `
cosmos -md 10 -c equinox-test-archive -a equinox-test-aux `
cosmos -c equinox-test
```

<a name="sqlstreamstore"></a>
9. Use [SqlStreamStore](https://github.com/SQLStreamStore/SQLStreamStore)
11. Use [SqlStreamStore](https://github.com/SQLStreamStore/SQLStreamStore)

The SqlStreamStore consists of:

Expand Down Expand Up @@ -461,8 +503,11 @@ For EventStore, the tests assume a running local instance configured as follows

### Provisioning CosmosDb (when not using -sc)

dotnet run -f netcoreapp3.1 -p tools/Equinox.Tool -- init -ru 400 `
dotnet run -p tools/Equinox.Tool -- init -ru 400 `
cosmos -s $env:EQUINOX_COSMOS_CONNECTION -d $env:EQUINOX_COSMOS_DATABASE -c $env:EQUINOX_COSMOS_CONTAINER
# Same for a Secondary container for integration testing of the fallback mechanism
dotnet run -p tools/Equinox.Tool -- init -ru 400 `
cosmos -s $env:EQUINOX_COSMOS_CONNECTION -d $env:EQUINOX_COSMOS_DATABASE -c $env:EQUINOX_COSMOS_CONTAINER2

### Provisioning SqlStreamStore

Expand Down Expand Up @@ -500,7 +545,8 @@ All non-alpha releases derive from tagged commits on `master`. The tag defines t

- [Provision](#provisioning):
- Start Local EventStore running in simulated cluster mode
- Set Environment variables X 3 for a CosmosDb database and container (you might need to `eqx init`)
- Set environment variables x 4 for a CosmosDB database and container (you might need to `eqx init`)
- Add a `EQUINOX_COSMOS_CONTAINER2` environment variable referencing a separate (`eqx init` initialized) CosmosDB Container that will be used to store fallback events in the [Fallback mechanism's tests](https://github.com/jet/equinox/pull/247)
- `docker-compose up` to start 3 servers for the `SqlStreamStore.*.Integration` test suites
- [NB `SqlStreamStore.MsSql` has not been tested yet](https://github.com/jet/equinox/issues/175) :see_no_evil: **
- Run `./build.ps1` in PowerShell (or PowerShell Core on MacOS via `brew install cask pwsh`)
Expand Down
41 changes: 32 additions & 9 deletions samples/Infrastructure/Storage.fs
Original file line number Diff line number Diff line change
Expand Up @@ -42,6 +42,9 @@ module Cosmos =
| [<AltCommandLine "-s">] Connection of string
| [<AltCommandLine "-d">] Database of string
| [<AltCommandLine "-c">] Container of string
| [<AltCommandLine "-s2">] Connection2 of string
| [<AltCommandLine "-d2">] Database2 of string
| [<AltCommandLine "-c2">] Container2 of string
interface IArgParserTemplate with
member a.Usage =
match a with
Expand All @@ -53,11 +56,20 @@ module Cosmos =
| Connection _ -> "specify a connection string for a Cosmos account. (optional if environment variable EQUINOX_COSMOS_CONNECTION specified)"
| Database _ -> "specify a database name for store. (optional if environment variable EQUINOX_COSMOS_DATABASE specified)"
| Container _ -> "specify a container name for store. (optional if environment variable EQUINOX_COSMOS_CONTAINER specified)"
| Connection2 _ -> "specify a connection string for Secondary Cosmos account. Default: use same as Primary Connection"
| Database2 _ -> "specify a database name for Secondary store. Default: use same as Primary Database"
| Container2 _ -> "specify a container name for store. Default: use same as Primary Container"
type Info(args : ParseResults<Arguments>) =
member __.Mode = args.GetResult(ConnectionMode,Microsoft.Azure.Cosmos.ConnectionMode.Direct)
member __.Connection = args.TryGetResult Connection |> defaultWithEnvVar "EQUINOX_COSMOS_CONNECTION" "Connection"
member __.Database = args.TryGetResult Database |> defaultWithEnvVar "EQUINOX_COSMOS_DATABASE" "Database"
member __.Container = args.TryGetResult Container |> defaultWithEnvVar "EQUINOX_COSMOS_CONTAINER" "Container"
member private __.Connection2 = args.TryGetResult Connection2
member private __.Database2 = args.TryGetResult Database2 |> Option.defaultWith (fun () -> __.Database)
member private __.Container2 = args.TryGetResult Container2 |> Option.defaultWith (fun () -> __.Container)
member __.Secondary = if args.Contains Connection2 || args.Contains Database2 || args.Contains Container2
then Some (__.Connection2, __.Database2, __.Container2)
else None

member __.Timeout = args.GetResult(Timeout,5.) |> TimeSpan.FromSeconds
member __.Retries = args.GetResult(Retries,1)
Expand All @@ -70,17 +82,28 @@ module Cosmos =
open Equinox.CosmosStore
open Serilog

let logContainer (log: ILogger) name (mode, endpoint, db, container) =
log.Information("CosmosDb {name:l} {mode} {connection} Database {database} Container {container}", name, mode, endpoint, db, container)
let connect (a : Info) conn =
let discovery = Discovery.ConnectionString conn
CosmosStoreClientFactory(a.Timeout, a.Retries, a.MaxRetryWaitTime, mode=a.Mode).Create(discovery)
let conn (log: ILogger) (a : Info) =
let discovery = Discovery.ConnectionString a.Connection
let client = CosmosStoreClientFactory(a.Timeout, a.Retries, a.MaxRetryWaitTime, mode=a.Mode).Create(discovery)
log.Information("CosmosDb {mode} {connection} Database {database} Container {container}",
a.Mode, client.Endpoint, a.Database, a.Container)
log.Information("CosmosDb timeout {timeout}s; Throttling retries {retries}, max wait {maxRetryWaitTime}s",
(let t = a.Timeout in t.TotalSeconds), a.Retries, let x = a.MaxRetryWaitTime in x.TotalSeconds)
client, a.Database, a.Container
let (primaryClient, primaryDatabase, primaryContainer) as primary = connect a a.Connection, a.Database, a.Container
logContainer log "Primary" (a.Mode, primaryClient.Endpoint, primaryDatabase, primaryContainer)
let secondary =
match a.Secondary with
| Some (Some c2, db, container) -> Some (connect a c2, db, container)
| Some (None, db, container) -> Some (primaryClient, db, container)
| None -> None
secondary |> Option.iter (fun (client, db, container) -> logContainer log "Secondary" (a.Mode, client.Endpoint, db, container))
primary, secondary
let config (log: ILogger) (cache, unfolds, batchSize) info =
let client, databaseId, containerId = conn log info
let conn = CosmosStoreConnection(client, databaseId, containerId)
let conn =
match conn log info with
| (client, databaseId, containerId), None ->
CosmosStoreConnection(client, databaseId, containerId)
| (client, databaseId, containerId), Some (client2, db2, cont2) ->
CosmosStoreConnection(client, databaseId, containerId, client2=client2, databaseId2=db2, containerId2=cont2)
let ctx = CosmosStoreContext(conn, defaultMaxItems = batchSize)
let cacheStrategy = match cache with Some c -> CachingStrategy.SlidingWindow (c, TimeSpan.FromMinutes 20.) | None -> CachingStrategy.NoCaching
StorageConfig.Cosmos (ctx, cacheStrategy, unfolds)
Expand Down
Loading