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
5 changes: 4 additions & 1 deletion CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,9 @@ The `Unreleased` section name is replaced by the expected version of next releas
## [Unreleased]

### Added

- `Cosmos`: `Prune` API to delete events from the head of a stream [#233](https://github.com/jet/equinox/pull/233)

### Changed
### Removed
### Fixed
Expand All @@ -29,7 +32,7 @@ The `Unreleased` section name is replaced by the expected version of next releas

### Added

- Add `eqx dump -b`, enabling overriding of Max Events per Batch
- `eqx dump -b`, enabling overriding of Max Events per Batch
- `MemoryStore`: Add `Committed` event to enable simulating Change Feeds in integration tests re [#205](https://github.com/jet/equinox/issues/205) [#221](https://github.com/jet/equinox/pull/221)

### Changed
Expand Down
5 changes: 4 additions & 1 deletion samples/Store/Integration/LogIntegration.fs
Original file line number Diff line number Diff line change
Expand Up @@ -40,6 +40,9 @@ module EquinoxCosmosInterop =
| Log.SyncSuccess m -> "CosmosSync200", m, None, m.ru
| Log.SyncConflict m -> "CosmosSync409", m, None, m.ru
| Log.SyncResync m -> "CosmosSyncResync", m, None, m.ru
| Log.PruneResponse m -> "CosmosPruneResponse", m, None, m.ru
| Log.Delete m -> "CosmosDelete", m, None, m.ru
| Log.Prune (events, m) -> "CosmosPrune", m, Some events, m.ru
{ action = action; stream = metric.stream; bytes = metric.bytes; count = metric.count; responses = batches
interval = StopwatchInterval(metric.interval.StartTicks,metric.interval.EndTicks); ru = ru }

Expand Down Expand Up @@ -128,4 +131,4 @@ type Tests() =
let itemCount = batchSize / 2 + 1
let cartId = % Guid.NewGuid()
do! act buffer service itemCount context cartId skuId "EqxCosmos Tip " // one is a 404, one is a 200
}
}
148 changes: 145 additions & 3 deletions src/Equinox.Cosmos/Cosmos.fs
Original file line number Diff line number Diff line change
Expand Up @@ -226,10 +226,20 @@ module Log =
/// Summarizes a set of Responses for a given Read request
| Query of Direction * responses: int * Measurement
/// Individual read request in a Batch
/// Charges are rolled up into Query (so do not double count)
| Response of Direction * Measurement
| SyncSuccess of Measurement
| SyncResync of Measurement
| SyncConflict of Measurement
/// Handled response from listing of batches in a stream
/// Charges are rolled up into Prune (so do not double count)
| PruneResponse of Measurement
/// Deleted an individual Batch
| Delete of Measurement
/// Pruned batches from head of a stream
/// Count in Measurement is number of batches (documents) deleted
/// Bytes in Measurement is number of events deleted
| Prune of responsesHandled : int * Measurement
let prop name value (log : ILogger) = log.ForContext(name, value)
let propData name (events: #IEventData<byte[]> seq) (log : ILogger) =
let render = function null -> "null" | bytes -> System.Text.Encoding.UTF8.GetString bytes
Expand Down Expand Up @@ -265,16 +275,20 @@ module Log =
module Stats =
let inline (|Stats|) ({ interval = i; ru = ru }: Measurement) = ru, let e = i.Elapsed in int64 e.TotalMilliseconds

let (|CosmosReadRc|CosmosWriteRc|CosmosResyncRc|CosmosResponseRc|) = function
let (|CosmosReadRc|CosmosWriteRc|CosmosResyncRc|CosmosResponseRc|CosmosDeleteRc|CosmosPruneRc|) = function
| Tip (Stats s)
| TipNotFound (Stats s)
| TipNotModified (Stats s)
| Query (_,_, (Stats s)) -> CosmosReadRc s
// slices are rolled up into batches so be sure not to double-count
| Response (_,(Stats s)) -> CosmosResponseRc s
| Response (_,(Stats s))
// costs roll up into Prune operation so be sure not to double-count
| PruneResponse (Stats s) -> CosmosResponseRc s
| SyncSuccess (Stats s)
| SyncConflict (Stats s) -> CosmosWriteRc s
| SyncResync (Stats s) -> CosmosResyncRc s
| Delete (Stats s) -> CosmosDeleteRc s
| Prune (_, (Stats s)) -> CosmosPruneRc s
let (|SerilogScalar|_|) : LogEventPropertyValue -> obj option = function
| (:? ScalarValue as x) -> Some x.Value
| _ -> None
Expand All @@ -294,10 +308,14 @@ module Log =
static member val Read = Counter.Create() with get, set
static member val Write = Counter.Create() with get, set
static member val Resync = Counter.Create() with get, set
static member val Delete = Counter.Create() with get, set
static member val Prune = Counter.Create() with get, set
static member Restart() =
LogSink.Read <- Counter.Create()
LogSink.Write <- Counter.Create()
LogSink.Resync <- Counter.Create()
LogSink.Delete <- Counter.Create()
LogSink.Prune <- Counter.Create()
let span = epoch.Elapsed
epoch.Restart()
span
Expand All @@ -306,6 +324,9 @@ module Log =
| CosmosMetric (CosmosReadRc stats) -> LogSink.Read.Ingest stats
| CosmosMetric (CosmosWriteRc stats) -> LogSink.Write.Ingest stats
| CosmosMetric (CosmosResyncRc stats) -> LogSink.Resync.Ingest stats
| CosmosMetric (CosmosDeleteRc stats) -> LogSink.Delete.Ingest stats
| CosmosMetric (CosmosPruneRc stats) -> LogSink.Prune.Ingest stats
| CosmosMetric (CosmosResponseRc _) -> () // Costs are already included in others
| _ -> ()

/// Relies on feeding of metrics from Log through to Stats.LogSink
Expand All @@ -314,7 +335,9 @@ module Log =
let stats =
[ "Read", Stats.LogSink.Read
"Write", Stats.LogSink.Write
"Resync", Stats.LogSink.Resync ]
"Resync", Stats.LogSink.Resync
"Delete", Stats.LogSink.Delete
"Prune", Stats.LogSink.Prune ]
let mutable rows, totalCount, totalRc, totalMs = 0, 0L, 0., 0L
let logActivity name count rc lat =
log.Information("{name}: {count:n0} requests costing {ru:n0} RU (average: {avg:n2}); Average latency: {lat:n0}ms",
Expand Down Expand Up @@ -762,6 +785,114 @@ module internal Tip =
let t = StopwatchInterval(startTicks, endTicks)
log |> logQuery direction maxItems stream t (!responseCount,allSlices.ToArray()) -1L ru }

// Manages deletion of batches
// Note: it's critical that we delete individually, in the correct order so as not to leave gaps
// Note: public so BatchIndices can be deserialized into
module Delete =

open FSharp.Control
open Microsoft.Azure.Documents.Linq

type BatchIndices = { id : string; i : int64; n : int64 }

let pruneBefore (log: ILogger) (container: Container, stream: string) maxItems beforePos : Async<int * int * int64> = async {
let! ct = Async.CancellationToken
let log = log |> Log.prop "stream" stream
let deleteItem id count : Async<float> = async {
let docLink = sprintf "%O/docs/%s" container.CollectionUri id
let qo = Client.RequestOptions(PartitionKey = PartitionKey stream)
let! t, res = container.Client.DeleteDocumentAsync(docLink, qo, cancellationToken=ct) |> Async.AwaitTaskCorrect |> Stopwatch.Time
let rc, ms = res.RequestCharge, (let e = t.Elapsed in e.TotalMilliseconds)
let reqMetric : Log.Measurement = { stream = stream; interval = t; bytes = -1; count = count; ru = rc }
let log = let evt = Log.Delete reqMetric in log |> Log.event evt
log.Information("EqxCosmos {action:l} {id} {ms}ms rc={ru}", "Delete", id, ms, rc)
return res.RequestCharge
}
let log = log |> Log.prop "beforePos" beforePos
let query : IDocumentQuery<BatchIndices> =
let q = SqlQuerySpec("SELECT c.id, c.i, c.n FROM c")
let qro = Client.FeedOptions(PartitionKey = PartitionKey stream, MaxItemCount=Nullable maxItems)
container.Client.CreateDocumentQuery<_>(container.CollectionUri, q, qro).AsDocumentQuery()
let tryReadNextPage (x : IDocumentQuery<_>) = async {
if not x.HasMoreResults then return None else

let! t, (res : Client.FeedResponse<_>) = query.ExecuteNextAsync<BatchIndices>(ct) |> Async.AwaitTaskCorrect |> Stopwatch.Time
let batches, rc, ms = Array.ofSeq res, res.RequestCharge, (let e = t.Elapsed in e.TotalMilliseconds)
let next = (match batches with [||] -> None | xs -> Some (xs.[ xs.Length - 1])) |> Option.map (fun x -> x.n) |> Option.toNullable
let reqMetric : Log.Measurement = { stream = stream; interval = t; bytes = -1; count = batches.Length; ru = rc }
let log = let evt = Log.PruneResponse reqMetric in log |> Log.event evt
log.Information("EqxCosmos {action:l} {batches} {ms}ms n={next} rc={ru}", "PruneResponse", batches.Length, ms, next, rc)
return Some ((rc, batches), x)
}
// If we have results: []
// - deleteBefore 9 would: return 0,0,0

// If we have results: ["-1",10,10; "0",0,1; "2",1,3; "3",3,8; "8",8,10]
// - deleteBefore 3 would: inspect first 3, delete 2, return 3,0,3
// - deleteBefore 4 would: inspect first 4, delete 2, return 3,1,3
// - deleteBefore 8 would: inspect first 4, delete 2, return 8,0,8

// If we have results: ["-1",10,10; "8",8,10]
// - deleteBefore 3 would: inspect first 2, delete 0, return 0,0,8
// - deleteBefore 8 would: inspect first 2, delete 0, return 0,0,8
// - deleteBefore 9 would: inspect first 2, delete 0, return 0,1,8
// - deleteBefore 10 would: inspect first 2, delete 1, return 2,0,10
// - deleteBefore 11 would: inspect first 2, delete 1, return 2,0,10

// If we have results: ["-1",10,10]
// - deleteBefore 9 would: inspect first 1, delete 0, return 0,0,10
// - deleteBefore 10 would: inspect first 1, delete 0, return 0,0,10
// - deleteBefore 11 would: inspect first 1, delete 0, return 0,0,10
let! pt, outcomes =
let isTip (x : BatchIndices) = x.id = Tip.WellKnownDocumentId
let isRelevant x = isTip x || x.i < beforePos
let hasRelevantItems (_, batches) = batches |> Array.exists isRelevant
let handle (rc, batches : BatchIndices[]) = async {
let mutable delCharges, batchesDeleted, eventsDeleted, eventsDeferred = 0., 0, 0, 0
let mutable tipI, lwm = None, None
for x in batches |> Seq.takeWhile isRelevant do
let count = x.n - x.i |> int
if isTip x then
tipI <- Some x.i
elif x.n > beforePos then
eventsDeferred <- eventsDeferred + min count (int (beforePos - x.i))
lwm <- Some x.i
else
let! charge = deleteItem x.id count
delCharges <- delCharges + charge
batchesDeleted <- batchesDeleted + 1
eventsDeleted <- eventsDeleted + count
lwm <- Some x.n
return rc, (tipI, lwm), (delCharges, batchesDeleted, eventsDeleted, eventsDeferred)
}
AsyncSeq.unfoldAsync tryReadNextPage query
|> AsyncSeq.takeWhile hasRelevantItems
|> AsyncSeq.mapAsync handle
|> AsyncSeq.toArrayAsync
|> Stopwatch.Time
let mutable queryCharges, delCharges, responses, batchesDeleted, eventsDeleted, eventsDeferred = 0., 0., 0, 0, 0, 0
let mutable lwm, tipI = None, None
for qc, (bTipI, bLwm), (dc, bDel, eDel, eDef) in outcomes do
lwm <- max lwm bLwm
tipI <- max tipI bTipI
queryCharges <- queryCharges + qc
delCharges <- delCharges + dc
responses <- responses + 1
batchesDeleted <- batchesDeleted + bDel
eventsDeleted <- eventsDeleted + eDel
eventsDeferred <- eventsDeferred + eDef
let reqMetric : Log.Measurement = { stream = stream; interval = pt; bytes = eventsDeleted; count = batchesDeleted; ru = queryCharges }
let log = let evt = Log.Prune (responses, reqMetric) in log |> Log.event evt
let lwm =
match lwm, tipI with
| Some lwm, _ -> lwm // we saw a batch and identified a Low Water mark based on it
| None, Some tipI -> tipI // we saw the Tip, but no batches along the way, therefore it's i is the low water mark
| None, None -> 0L // If we've seen no batches at all, then the write position is 0L
log.Information("EqxCosmos {action:l} {events}/{batches} lwm={lwm} {ms}ms queryRu={ru} deleteRu={delRu}",
"Prune", eventsDeleted, batchesDeleted, lwm, (let e = pt.Elapsed in e.TotalMilliseconds), queryCharges, delCharges)
return eventsDeleted, eventsDeferred, lwm
}

type [<NoComparison>] Token = { container: Container; stream: string; pos: Position }
module Token =
let create (container,stream) pos : StreamToken =
Expand Down Expand Up @@ -866,6 +997,8 @@ type Gateway(conn : Connection, batching : BatchingPolicy) =
| Sync.Result.Conflict (pos',events) -> return InternalSyncResult.Conflict (Token.create containerStream pos',events)
| Sync.Result.ConflictUnknown pos' -> return InternalSyncResult.ConflictUnknown (Token.create containerStream pos')
| Sync.Result.Written pos' -> return InternalSyncResult.Written (Token.create containerStream pos') }
member __.Prune(log, (container, stream), beforeIndex) =
Delete.pruneBefore log (container, stream) batching.MaxItems beforeIndex

type private Category<'event, 'state, 'context>(gateway : Gateway, codec : IEventCodec<'event,byte[],'context>) =
let (|TryDecodeFold|) (fold: 'state -> 'event seq -> 'state) initial (events: ITimelineEvent<byte[]> seq) : 'state = Seq.choose codec.TryDecode events |> fold initial
Expand Down Expand Up @@ -1280,6 +1413,9 @@ type Context
| AppendResult.Ok token -> return token
| x -> return x |> sprintf "Conflict despite it being disabled %A" |> invalidOp }

member __.Prune((container,stream), beforeIndex) : Async<int * int * int64> =
gateway.Prune(log, (container,stream), beforeIndex)

/// Provides mechanisms for building `EventData` records to be supplied to the `Events` API
type EventData() =
/// Creates an Event record, suitable for supplying to Append et al
Expand Down Expand Up @@ -1335,6 +1471,12 @@ module Events =
let appendAtEnd (ctx: Context) (streamName: string) (events: IEventData<_>[]): Async<int64> =
ctx.NonIdempotentAppend(ctx.CreateStream streamName, events) |> stripPosition

/// Requests deletion of events prior to the specified Index
/// Due to the need to preserve ordering of data in the stream, only full batches will be removed
/// Returns count of events deleted this time, events that could not be deleted due to partial batches, and the stream's lowest remaining sequence number
let prune (ctx: Context) (streamName: string) (beforeIndex: int64): Async<int * int * int64> =
ctx.Prune(ctx.CreateStream streamName, beforeIndex)

/// Returns an async sequence of events in the stream backwards starting from the specified sequence number,
/// reading in batches of the specified size.
/// Returns an empty sequence if the stream is empty or if the sequence number is smaller than the smallest
Expand Down
47 changes: 47 additions & 0 deletions tests/Equinox.Cosmos.Integration/CosmosCoreIntegration.fs
Original file line number Diff line number Diff line change
Expand Up @@ -325,3 +325,50 @@ type Tests(testOutputHelper) =
[1,5] =! capture.ChooseCalls queryRoundTripsAndItemCounts
verifyRequestChargesMax 4 // 3.24 // WAS 3 // 2.98
}

(* Prune *)
[<AutoData(SkipIfRequestedViaEnvironmentVariable="EQUINOX_INTEGRATION_SKIP_COSMOS")>]
let prune (TestStream streamName) = Async.RunSynchronously <| async {
capture.Clear()
let! conn = connectToSpecifiedCosmosOrSimulator log
let ctx = mkContextWithItemLimit conn None

let! expected = add6EventsIn2Batches ctx streamName

// Trigger deletion of first batch
capture.Clear()
let! deleted, deferred, trimmedPos = Events.prune ctx streamName 5L
test <@ deleted = 1 && deferred = 4 && trimmedPos = 1L @>
test <@ [EqxAct.PruneResponse; EqxAct.Delete; EqxAct.Prune] = capture.ExternalCalls @>
verifyRequestChargesMax 17 // 13.33 + 2.9

let! res = Events.get ctx streamName 0L Int32.MaxValue
verifyCorrectEvents 1L (Array.skip 1 expected) res

// Repeat the process, but this time there should be no actual deletes
capture.Clear()
let! deleted, deferred, trimmedPos = Events.prune ctx streamName 4L
test <@ deleted = 0 && deferred = 3 && trimmedPos = 1L @>
test <@ [EqxAct.PruneResponse; EqxAct.Prune] = capture.ExternalCalls @>
verifyRequestChargesMax 3 // 2.86

let! res = Events.get ctx streamName 0L Int32.MaxValue
verifyCorrectEvents 1L (Array.skip 1 expected) res

// Delete second batch
capture.Clear()
let! deleted, deferred, trimmedPos = Events.prune ctx streamName 6L
test <@ deleted = 5 && deferred = 0 && trimmedPos = 6L @>
test <@ [EqxAct.PruneResponse; EqxAct.Delete; EqxAct.Prune] = capture.ExternalCalls @>
verifyRequestChargesMax 17 // 13.33 + 2.86

let! res = Events.get ctx streamName 0L Int32.MaxValue
test <@ [||] = res @>

// Attempt to repeat
capture.Clear()
let! deleted, deferred, trimmedPos = Events.prune ctx streamName 6L
test <@ deleted = 0 && deferred = 0 && trimmedPos = 6L @>
test <@ [EqxAct.PruneResponse; EqxAct.Prune] = capture.ExternalCalls @>
verifyRequestChargesMax 3 // 2.83
}
Loading