Test interleaved reading/writing
| 485 | } |
| 486 | // Test interleaved reading/writing |
| 487 | void DataTest::TestDoExchangeTotal() { |
| 488 | auto descr = FlightDescriptor::Command("total"); |
| 489 | std::unique_ptr<FlightStreamReader> reader; |
| 490 | std::unique_ptr<FlightStreamWriter> writer; |
| 491 | { |
| 492 | auto a1 = ArrayFromJSON(arrow::int32(), "[4, 5, 6, null]"); |
| 493 | auto schema = arrow::schema({field("f1", a1->type())}); |
| 494 | // XXX: as noted in flight/client.cc, Begin() is lazy and the |
| 495 | // schema message won't be written until some data is also |
| 496 | // written. There's also timing issues; hence we check each status |
| 497 | // here. |
| 498 | EXPECT_RAISES_WITH_MESSAGE_THAT( |
| 499 | Invalid, ::testing::HasSubstr("Field is not INT64: f1"), ([&]() { |
| 500 | ARROW_ASSIGN_OR_RAISE(auto exchange, client_->DoExchange(descr)); |
| 501 | reader = std::move(exchange.reader); |
| 502 | writer = std::move(exchange.writer); |
| 503 | RETURN_NOT_OK(writer->Begin(schema)); |
| 504 | auto batch = RecordBatch::Make(schema, /* num_rows */ 4, {a1}); |
| 505 | RETURN_NOT_OK(writer->WriteRecordBatch(*batch)); |
| 506 | return writer->Close(); |
| 507 | })()); |
| 508 | } |
| 509 | { |
| 510 | auto a1 = ArrayFromJSON(arrow::int64(), "[1, 2, null, 3]"); |
| 511 | auto a2 = ArrayFromJSON(arrow::int64(), "[null, 4, 5, 6]"); |
| 512 | auto schema = arrow::schema({field("f1", a1->type()), field("f2", a2->type())}); |
| 513 | ASSERT_OK_AND_ASSIGN(auto exchange, client_->DoExchange(descr)); |
| 514 | reader = std::move(exchange.reader); |
| 515 | writer = std::move(exchange.writer); |
| 516 | ASSERT_OK(writer->Begin(schema)); |
| 517 | auto batch = RecordBatch::Make(schema, /* num_rows */ 4, {a1, a2}); |
| 518 | ASSERT_OK(writer->WriteRecordBatch(*batch)); |
| 519 | ASSERT_OK_AND_ASSIGN(auto server_schema, reader->GetSchema()); |
| 520 | AssertSchemaEqual(*schema, *server_schema); |
| 521 | |
| 522 | ASSERT_OK_AND_ASSIGN(auto chunk, reader->Next()); |
| 523 | ASSERT_NE(nullptr, chunk.data); |
| 524 | auto expected1 = RecordBatch::Make( |
| 525 | schema, /* num_rows */ 1, |
| 526 | {ArrayFromJSON(arrow::int64(), "[6]"), ArrayFromJSON(arrow::int64(), "[15]")}); |
| 527 | AssertBatchesEqual(*expected1, *chunk.data); |
| 528 | |
| 529 | ASSERT_OK(writer->WriteRecordBatch(*batch)); |
| 530 | ASSERT_OK_AND_ASSIGN(chunk, reader->Next()); |
| 531 | ASSERT_NE(nullptr, chunk.data); |
| 532 | auto expected2 = RecordBatch::Make( |
| 533 | schema, /* num_rows */ 1, |
| 534 | {ArrayFromJSON(arrow::int64(), "[12]"), ArrayFromJSON(arrow::int64(), "[30]")}); |
| 535 | AssertBatchesEqual(*expected2, *chunk.data); |
| 536 | |
| 537 | ASSERT_OK(writer->Close()); |
| 538 | } |
| 539 | } |
| 540 | // Ensure server errors get propagated no matter what we try |
| 541 | void DataTest::TestDoExchangeError() { |
| 542 | auto descr = FlightDescriptor::Command("error"); |
nothing calls this directly
no test coverage detected