-
-
Notifications
You must be signed in to change notification settings - Fork 2
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
On error resume #2
Comments
This comment has been minimized.
This comment has been minimized.
Other option would be to simply use OnNext instead of OnError for connection issues and wrap inside an F# result module temp.RedisRx
open System
open System.Threading
open System.Reactive.Linq
open System.Reactive
open System.Threading.Tasks
open FSharp.Control.Reactive
open FSharp.Control.Redis.Streams.Core
open FSharp.Control.Tasks
open StackExchange.Redis
module Observable =
let internal taskUnfold (fn: 's -> CancellationToken -> Task<('s * 'e) option>) (state: 's) =
Observable.Create(fun (obs: IObserver<_>) ->
let cts = new CancellationTokenSource()
let ct = cts.Token
task {
let mutable innerState = state
let mutable isFinished = false
try
while not ct.IsCancellationRequested || not isFinished do
try
let! result = fn innerState ct
match result with
| Some (newState, output) ->
innerState <- newState
obs.OnNext <| Ok output
| None ->
isFinished <- true
with e ->
obs.OnNext <| Error e
obs.OnCompleted()
finally
cts.Dispose()
}
|> ignore
new Disposables.CancellationDisposable(cts) :> IDisposable)
let flattenArray (observable : IObservable<Result<StreamEntry [],exn>>) =
observable.SelectMany (fun x ->
match x with
| Ok o -> o |> Array.map (Ok) :> seq<_>
| Error e -> [Error e] :> seq<_>)
let pollStreamForever (redisdb : IDatabase) (streamName : RedisKey) (startingPosition : RedisValue) (pollOptions : PollOptions) =
Observable.taskUnfold (fun (nextPosition, pollDelay) ct -> task {
let! (response : StreamEntry []) = redisdb.StreamRangeAsync(streamName, minId = Nullable(nextPosition), count = (Option.toNullable pollOptions.CountToPullATime))
match response with
| EmptyArray ->
let nextPollDelay = pollOptions.CalculateNextPollDelay pollDelay
do! Task.Delay pollDelay
return Some ((nextPosition, nextPollDelay ) , Array.empty )
| entries ->
let lastEntry = Seq.last entries
let nextPosition = EntryId.CalculateNextPositionIncr lastEntry.Id
let nextPollDelay = TimeSpan.Zero
return Some ((nextPosition, nextPollDelay), entries )
}) (startingPosition, TimeSpan.Zero)
|> Observable.flattenArray
|
We discussed this in slack, creating a new function like Could you open a PR for this, also adding this function for the |
Slack convo:
|
6 tasks
@TheAngryByrd I have created the pull request. 👍 |
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Currently using pollStreamForever like,
Is it possible to resume where it left off on Error?
The text was updated successfully, but these errors were encountered: