NSUKit 1.4.0
板卡级统一交互接口
载入中...
搜索中...
未找到
data_download.cpp
浏览该文件的文档.
1
2// Copyright (c) 2024. Naishu
3// NSUKit is licensed under Mulan PSL v2.
4// You can use this software according to the terms and conditions of the Mulan PSL v2.
5// You may obtain a copy of Mulan PSL v2 at:
6// http://license.coscl.org.cn/MulanPSL2
7// THIS SOFTWARE IS PROVIDED ON AN "AS IS" BASIS, WITHOUT WARRANTIES OF ANY KIND,
8// EITHER EXPRESS OR IMPLIED, INCLUDING BUT NOT LIMITED TO NON-INFRINGEMENT,
9// MERCHANTABILITY OR FIT FOR A PARTICULAR PURPOSE.
10// See the Mulan PSL v2 for more details.
12
13//
14// Created by jilianyi<jilianyi@naishu.tech> on 2024/4/12.
15//
16
17/*
18 * +-----------------+ +-------------------+
19 * | | +---------------+ | |
20 * | |-----> | full_queue |---- ptr -->| |
21 * | | +---------------+ | |
22 * | download_thread | | data_generate |
23 * | | +---------------+ | |
24 * | |<----- | empty_queue |<--- ptr ---| |
25 * | | +---------------+ | |
26 * +-----------------+ +-------------------+
27 * @class ThreadSafeQueue
28 */
29#include <iostream>
30#include <queue>
31#include <thread>
32#include <condition_variable>
33#include "NSUKit.h"
34
35
40template <typename T>
42public:
44
45 //
46 void Push(T *value) {
47 std::unique_lock<std::mutex> lock(mutex_);
48 queue_.push(value);
49 lock.unlock();
50 condition_.notify_one();
51 }
52
53 //
54 T *Pop() {
55 std::unique_lock<std::mutex> lock(mutex_);
56 condition_.wait(lock, [this] { return !queue_.empty(); });
57 auto value = (T *)queue_.front();
58 queue_.pop();
59 return value;
60 }
61
62private:
63 std::queue<T *> queue_;
64 std::mutex mutex_;
65 std::condition_variable condition_;
66};
67
68
70
71
72struct Deque {
75 int stopFlag = 0;
76};
77
78
87void download_thread(nsukit::BaseKit *_kit, Deque *q, nsuSize_t block, nsuSize_t total, std::mutex *mu) {
88 std::unique_lock<std::mutex> *lock;
89 nsuMemory_p mem;
90 nsuSize_t current = 0;
91 while (true) {
92 mem = q->empty.Pop();
93 auto s = _kit->open_send(0, mem, block, 0);
95 std::cout << "Establish CS and CR connections: " << std::endl;
96 std::cout << "Failed to enable data uplink " << nsukit::status2_string(s) << ", currently uplinked: " << current << std::endl;
97 break;
98 }
101 s = _kit->wait_stream(mem, 1.);
102 }
103 current += block;
104 q->full.Push(mem);
105 if (current == total) break;
106 }
107 lock = new std::unique_lock<std::mutex>(*mu);
108 q->stopFlag = 1;
109 delete lock;
110}
111
112
120void data_generate(nsukit::BaseKit *_kit, Deque *q, nsuSize_t block, std::mutex *mu) {
121 std::unique_lock<std::mutex> *lock;
122 nsuMemory_p mem;
123 unsigned int *buffer;
124 int cnt = 0;
125 nsuSize_t speed_count;
126 auto st = std::chrono::steady_clock::now();
127
128 while (true) {
129 lock = new std::unique_lock<std::mutex>(*mu);
130 if (q->stopFlag == 1) {
131 break;
132 }
133 delete lock;
134 mem = q->full.Pop();
135 buffer = (unsigned int *)_kit->get_buffer(mem, block);
136 for(int i=0; i<block; i++) buffer[i] = i;
137 q->empty.Push(mem);
138 speed_count += block;
139
140 if (cnt % 20 == 0) {
141 auto count = std::chrono::steady_clock::now() - st;
142 std::cout << std::flush << '\r' << "当前数据生成速度: " << speed_count*1000./(count.count()) << "MB/s" << std::endl;
143// speed_count = 0;
144 }
145 }
146 std::cout << std::endl;
147}
148
149
150int main(int argc, char *argv[]) {
151 unsigned int ds_block = 1024*1024;
152 Deque q;
153 std::mutex mu;
154 nsukit::NSUSoc <nsukit::TCPCmdUItf, nsukit::PCIECmdUItf, nsukit::PCIEStreamUItf> kit{};
155
156 if (argc != 4) {
157 std::cout << "Unsupported parameter passing method" << std::endl;
158 // DataUpload 127.0.0.1 104857600 ./da_data_1.dat
159 std::cout << argv[0] << " {IP} {totalBytes} {filePath}" << std::endl;
160 return 1;
161 }
162 nsuSize_t total_len = std::atoi(argv[2]);
163 if (total_len % ds_block != 0) {
164 std::cout << "The total length of upstream data total_len "
165 << total_len << " Bytes should be "
166 << ds_block << " Integer multiple of Bytes" << std::endl;
167 return 1;
168 }
169
170 nsuInitParam_t param;
171 param.cmd_ip = argv[1];
172 param.cmd_board = 0;
173 param.stream_board = 0;
174
175 auto res = kit.link_cmd(&param);
177 std::cout << "Establish CS and CR connections: " << nsukit::status2_string(res) << std::endl;
178 }
179
180 res = kit.link_stream(&param);
182 std::cout << "Establish a DS connection: " << nsukit::status2_string(res) << std::endl;
183 }
184
185 for (int i=0; i<10; i++) {
186 nsuMemory_p mem = kit.alloc_buffer(ds_block);
187 q.empty.Push(mem);
188 }
189
190
191 // 通知FPGA开始采集
192 res = kit.execute("系统开启");
194 std::cout << "系统开启:" << nsukit::status2_string(res) << std::endl;
195 }
196
197 // 开启上行及存盘线程
198 std::thread up_trd(download_thread, &kit, &q, ds_block, total_len, &mu);
199 std::thread generate_trd(data_generate, &kit, &q, ds_block, &mu);
200
201 up_trd.join();
202 generate_trd.join();
203
204 // 通知FPGA开始采集
205 res = kit.execute("系统停止");
207 std::cout << "系统停止:" << nsukit::status2_string(res) << std::endl;
208 }
209
210 std::cout << "Data download completed" << std::endl;
211
212 return 0;
213}
void Push(T *value)
virtual nsuVoidBuf_p get_buffer(nsuMemory_p fd, nsuStreamLen_t length=0)
Definition: base_kit.h:76
virtual nsukitStatus_t open_send(nsuChnlNum_t chnl, nsuMemory_p fd, nsuStreamLen_t length, nsuStreamLen_t offset=0)
Definition: base_kit.h:93
virtual nsukitStatus_t wait_stream(nsuMemory_p fd, float timeout=0.)
Definition: base_kit.h:103
void data_generate(nsukit::BaseKit *_kit, Deque *q, nsuSize_t block, std::mutex *mu)
void download_thread(nsukit::BaseKit *_kit, Deque *q, nsuSize_t block, nsuSize_t total, std::mutex *mu)
ThreadSafeQueue< void > memQueue
std::string NSU_DLLEXPORT status2_string(nsukitStatus_t status)
Definition: config.cpp:15
int stopFlag
memQueue empty
memQueue full
std::string cmd_ip
Definition: type.h:98
nsuBoardNum_t cmd_board
Definition: type.h:115
nsuBoardNum_t stream_board
Definition: type.h:121
@ NSUKIT_STATUS_STREAM_RUNNING
DLLEXTERN typedef size_t nsuSize_t
Definition: type.h:87
DLLEXTERN typedef void * nsuMemory_p
Definition: type.h:86