Shell: mempool RPC fixes
This commit is contained in:
parent
1357aff0ba
commit
d53918451b
@ -602,18 +602,19 @@ module Make(Static: STATIC)(Proto: Registered_protocol.T)
|
||||
let state = Worker.state w in
|
||||
let filter_result = function
|
||||
| Applied _ -> params#applied
|
||||
| Refused _ -> params#branch_refused
|
||||
| Branch_refused _ -> params#refused
|
||||
| Refused _ -> params#refused
|
||||
| Branch_refused _ -> params#branch_refused
|
||||
| Branch_delayed _ -> params#branch_delayed
|
||||
| _ -> false in
|
||||
|
||||
let op_stream, stopper = Lwt_watcher.create_stream state.operation_stream in
|
||||
let shutdown () = Lwt_watcher.shutdown stopper in
|
||||
let next () =
|
||||
let rec next () =
|
||||
Lwt_stream.get op_stream >>= function
|
||||
| Some (kind, shell, protocol_data) when filter_result kind ->
|
||||
Lwt.return_some [ { Proto.shell ; protocol_data } ]
|
||||
| _ -> Lwt.return_none in
|
||||
| Some _ -> next ()
|
||||
| None -> Lwt.return_none in
|
||||
RPC_answer.return_stream { next ; shutdown }
|
||||
)
|
||||
|
||||
|
@ -470,7 +470,7 @@ module Make(Proto: Registered_protocol.T)(Arg: ARG): T = struct
|
||||
(fun (hash, op) -> (hash, map_op op))
|
||||
(List.rev pv.applied) ;
|
||||
refused =
|
||||
Operation_hash.Map.map map_op_error pv.branch_refusals ;
|
||||
Operation_hash.Map.map map_op_error pv.refusals ;
|
||||
branch_refused =
|
||||
Operation_hash.Map.map map_op_error pv.branch_refusals ;
|
||||
branch_delayed =
|
||||
@ -510,11 +510,11 @@ module Make(Proto: Registered_protocol.T)(Arg: ARG): T = struct
|
||||
let current_mempool = ref (Some current_mempool) in
|
||||
let filter_result = function
|
||||
| `Applied -> params#applied
|
||||
| `Refused -> params#branch_refused
|
||||
| `Branch_refused -> params#refused
|
||||
| `Refused -> params#refused
|
||||
| `Branch_refused -> params#branch_refused
|
||||
| `Branch_delayed -> params#branch_delayed
|
||||
in
|
||||
let next () =
|
||||
let rec next () =
|
||||
match !current_mempool with
|
||||
| Some mempool -> begin
|
||||
current_mempool := None ;
|
||||
@ -532,7 +532,8 @@ module Make(Proto: Registered_protocol.T)(Arg: ARG): T = struct
|
||||
Proto.operation_data_encoding
|
||||
bytes in
|
||||
Lwt.return_some [ { Proto.shell ; protocol_data } ]
|
||||
| _ -> Lwt.return_none
|
||||
| Some _ -> next ()
|
||||
| None -> Lwt.return_none
|
||||
end
|
||||
in
|
||||
let shutdown () = Lwt_watcher.shutdown stopper in
|
||||
|
Loading…
Reference in New Issue
Block a user