|
| 1 | +// SPDX-License-Identifier: BSD-3-Clause |
| 2 | +/* Copyright 2021, Intel Corporation */ |
| 3 | + |
| 4 | +/** |
| 5 | + * mpsc_queue.cpp -- example which shows how to use |
| 6 | + * pmem::obj::experimental::mpsc_queue |
| 7 | + */ |
| 8 | + |
| 9 | +#include <libpmemobj++/experimental/mpsc_queue.hpp> |
| 10 | +#include <libpmemobj++/make_persistent.hpp> |
| 11 | +#include <libpmemobj++/persistent_ptr.hpp> |
| 12 | +#include <libpmemobj++/transaction.hpp> |
| 13 | + |
| 14 | +#include <iostream> |
| 15 | +#include <string> |
| 16 | + |
| 17 | +void |
| 18 | +show_usage(char *argv[]) |
| 19 | +{ |
| 20 | + std::cerr << "usage: " << argv[0] << " file-name" << std::endl; |
| 21 | +} |
| 22 | + |
| 23 | +//! [mpsc_queue_single_threaded_example] |
| 24 | + |
| 25 | +struct root { |
| 26 | + pmem::obj::persistent_ptr< |
| 27 | + pmem::obj::experimental::mpsc_queue::pmem_log_type> |
| 28 | + log; |
| 29 | +}; |
| 30 | + |
| 31 | +void |
| 32 | +single_threaded(pmem::obj::pool<root> pop) |
| 33 | +{ |
| 34 | + |
| 35 | + std::vector<std::string> values_to_produce = {"xxx", "aaaaaaa", "bbbbb", |
| 36 | + "cccc", "ddddddddddd"}; |
| 37 | + pmem::obj::persistent_ptr<root> proot = pop.root(); |
| 38 | + |
| 39 | + /* Create mpsc_queue, which uses pmem_log_type object to store |
| 40 | + * data. */ |
| 41 | + auto queue = pmem::obj::experimental::mpsc_queue(*proot->log, 1); |
| 42 | + |
| 43 | + /* Consume data, which was stored in the queue in the previous run of |
| 44 | + * the application. */ |
| 45 | + //! [try_consume_batch] |
| 46 | + queue.try_consume_batch( |
| 47 | + [&](pmem::obj::experimental::mpsc_queue::batch_type rd_acc) { |
| 48 | + for (pmem::obj::string_view str : rd_acc) { |
| 49 | + std::cout << std::string(str.data(), str.size()) |
| 50 | + << std::endl; |
| 51 | + } |
| 52 | + }); |
| 53 | + //! [try_consume_batch] |
| 54 | + /* Produce and consume data. */ |
| 55 | + //! [register_worker] |
| 56 | + pmem::obj::experimental::mpsc_queue::worker worker = |
| 57 | + queue.register_worker(); |
| 58 | + //! [register_worker] |
| 59 | + for (std::string &value : values_to_produce) { |
| 60 | + //! [try_produce] |
| 61 | + /* Produce data. */ |
| 62 | + worker.try_produce(value); |
| 63 | + //! [try_produce] |
| 64 | + /* Consume produced data. */ |
| 65 | + queue.try_consume_batch( |
| 66 | + [&](pmem::obj::experimental::mpsc_queue::batch_type |
| 67 | + rd_acc) { |
| 68 | + for (pmem::obj::string_view str : rd_acc) { |
| 69 | + std::cout << std::string(str.data(), |
| 70 | + str.size()) |
| 71 | + << std::endl; |
| 72 | + } |
| 73 | + }); |
| 74 | + } |
| 75 | + //! [try_produce_string_view] |
| 76 | + /* Produce data to be consumed in next run of the application. */ |
| 77 | + worker.try_produce("Left for next run"); |
| 78 | + //! [try_produce_string_view] |
| 79 | +} |
| 80 | + |
| 81 | +//! [mpsc_queue_single_threaded_example] |
| 82 | + |
| 83 | +int |
| 84 | +main(int argc, char *argv[]) |
| 85 | +{ |
| 86 | + if (argc < 2) { |
| 87 | + show_usage(argv); |
| 88 | + return 1; |
| 89 | + } |
| 90 | + |
| 91 | + const char *path = argv[1]; |
| 92 | + |
| 93 | + static constexpr size_t QUEUE_SIZE = 1000; |
| 94 | + |
| 95 | + pmem::obj::pool<root> pop; |
| 96 | + try { |
| 97 | + pop = pmem::obj::pool<root>::open(path, "mpsc_queue"); |
| 98 | + if (pop.root()->log == nullptr) { |
| 99 | + pmem::obj::transaction::run(pop, [&] { |
| 100 | + pop.root()->log = pmem::obj::make_persistent< |
| 101 | + pmem::obj::experimental::mpsc_queue:: |
| 102 | + pmem_log_type>(QUEUE_SIZE); |
| 103 | + }); |
| 104 | + } |
| 105 | + single_threaded(pop); |
| 106 | + |
| 107 | + } catch (pmem::pool_error &e) { |
| 108 | + std::cerr << e.what() << std::endl; |
| 109 | + std::cerr |
| 110 | + << "To create pool run: pmempool create obj --layout=mpsc_queue -s 100M path_to_pool" |
| 111 | + << std::endl; |
| 112 | + } catch (std::exception &e) { |
| 113 | + std::cerr << e.what() << std::endl; |
| 114 | + } |
| 115 | + |
| 116 | + try { |
| 117 | + pop.close(); |
| 118 | + } catch (const std::logic_error &e) { |
| 119 | + std::cerr << e.what() << std::endl; |
| 120 | + } |
| 121 | + return 0; |
| 122 | +} |
0 commit comments