(ctx context.Context, queryid, src string, backend sourcebackendpb.SourceBackendClient, backendidx int, searchRequest *sourcebackendpb.SearchRequest)
| 181 | ) |
| 182 | |
| 183 | func queryBackend(ctx context.Context, queryid, src string, backend sourcebackendpb.SourceBackendClient, backendidx int, searchRequest *sourcebackendpb.SearchRequest) { |
| 184 | // When exiting this function, check that all results were processed. If |
| 185 | // not, the backend query must have failed for some reason. Send a progress |
| 186 | // update to prevent the query from running forever. |
| 187 | defer func() { |
| 188 | stateMu.RLock() |
| 189 | s, ok := state[queryid] |
| 190 | stateMu.RUnlock() |
| 191 | if !ok { |
| 192 | return // query no longer exists |
| 193 | } |
| 194 | s.filesMu.Lock() |
| 195 | filesTotal := s.filesTotal[backendidx] |
| 196 | filesProcessed := s.filesProcessed[backendidx] |
| 197 | s.filesMu.Unlock() |
| 198 | |
| 199 | if filesProcessed == filesTotal { |
| 200 | return |
| 201 | } |
| 202 | |
| 203 | if filesTotal == -1 { |
| 204 | filesTotal = 0 |
| 205 | } |
| 206 | |
| 207 | storeProgress(queryid, backendidx, &sourcebackendpb.ProgressUpdate{ |
| 208 | FilesProcessed: uint64(filesTotal), |
| 209 | FilesTotal: uint64(filesTotal), |
| 210 | }) |
| 211 | |
| 212 | addEventMarshal(queryid, &Error{ |
| 213 | Type: "error", |
| 214 | ErrorType: "backendunavailable", |
| 215 | }) |
| 216 | }() |
| 217 | |
| 218 | ctx, cancelfunc := context.WithCancel(ctx) |
| 219 | stream, err := backend.Search(ctx, searchRequest) |
| 220 | if err != nil { |
| 221 | log.Printf("[%s] [src:%s] Search RPC failed: %v\n", queryid, src, err) |
| 222 | return |
| 223 | } |
| 224 | |
| 225 | stateMu.RLock() |
| 226 | s, ok := state[queryid] |
| 227 | stateMu.RUnlock() |
| 228 | if !ok { |
| 229 | log.Printf("[%s] [src:%s] query no longer exists\n", queryid, src) |
| 230 | return |
| 231 | } |
| 232 | bstate := s.perBackend[backendidx] |
| 233 | tempFileWriter := bstate.tempFileWriter |
| 234 | orderlyFinished := false |
| 235 | done := false |
| 236 | |
| 237 | for !done { |
| 238 | msg, err := stream.Recv() |
| 239 | if err == io.EOF { |
| 240 | log.Printf("[%s] [src:%s] EOF\n", queryid, src) |
no test coverage detected