(ctx context.Context, req *scatterRequest, pbReq *pb.SearchRequest)
| 242 | } |
| 243 | |
| 244 | func (g *GrpcClient) SearchRPC(ctx context.Context, req *scatterRequest, |
| 245 | pbReq *pb.SearchRequest) (*bleve.SearchResult, error) { |
| 246 | res, err := g.GrpcCli.Search(ctx, pbReq) |
| 247 | if err != nil || res == nil { |
| 248 | log.Errorf("grpc_client: search err: %v", err) |
| 249 | return nil, err |
| 250 | } |
| 251 | |
| 252 | searchResult := &bleve.SearchResult{ |
| 253 | Status: &bleve.SearchStatus{ |
| 254 | Errors: make(map[string]error)}, |
| 255 | Request: req.searchRequest, |
| 256 | } |
| 257 | |
| 258 | var response *pb.StreamSearchResults |
| 259 | for { |
| 260 | response, err = res.Recv() |
| 261 | if err == io.EOF { |
| 262 | err = nil |
| 263 | break |
| 264 | } |
| 265 | if err != nil { |
| 266 | break |
| 267 | } |
| 268 | |
| 269 | switch r := response.Contents.(type) { |
| 270 | |
| 271 | case *pb.StreamSearchResults_Hits: |
| 272 | if sw, ok := g.sc.(streamHandler); ok { |
| 273 | err = sw.write(r.Hits.Bytes, r.Hits.Offsets, int(r.Hits.Total)) |
| 274 | if err != nil { |
| 275 | break |
| 276 | } |
| 277 | } |
| 278 | |
| 279 | case *pb.StreamSearchResults_SearchResult: |
| 280 | if r.SearchResult != nil { |
| 281 | err = UnmarshalJSON(r.SearchResult, &searchResult) |
| 282 | if err != nil { |
| 283 | return searchResult, err |
| 284 | } |
| 285 | } |
| 286 | } |
| 287 | } |
| 288 | |
| 289 | return searchResult, err |
| 290 | } |
| 291 | |
| 292 | func (g *GrpcClient) Query(ctx context.Context, |
| 293 | req *scatterRequest) (*bleve.SearchResult, error) { |
no test coverage detected