33#include <condition_variable>
37#define STREAM_WITH_PCIE
39#define SocType nsukit::NSUSoc <nsukit::SimCmdUItf, nsukit::SimCmdUItf, nsukit::SimStreamUItf>
41#ifdef STREAM_WITH_PCIE
42#define SocType nsukit::NSUSoc <nsukit::PCIECmdUItf, nsukit::PCIECmdUItf, nsukit::PCIEStreamUItf>
46#define SocType nsukit::NSUSoc <nsukit::TCPCmdUItf, nsukit::TCPCmdUItf, nsukit::TCPStreamUItf>
61 std::unique_lock<std::mutex> lock(mutex_);
64 condition_.notify_one();
69 std::unique_lock<std::mutex> lock(mutex_);
70 condition_.wait(lock, [
this] {
return !queue_.empty(); });
71 auto value = (T *)queue_.front();
77 std::queue<T *> queue_;
79 std::condition_variable condition_;
102 std::unique_lock<std::mutex> *lock;
107 auto s = _kit->
open_recv(stream_chnl, mem, block, 0);
109 std::cout <<
"Establish CS and CR connections: " << std::endl;
110 std::cout <<
"Failed to enable data uplink " <<
nsukit::status2_string(s) <<
", currently uplinked: " << current << std::endl;
119 if (current == total) {
123 lock =
new std::unique_lock<std::mutex>(*mu);
138 std::unique_lock<std::mutex> *lock;
141 outF.open(path, std::ofstream::binary);
144 auto st = std::chrono::steady_clock::now();
147 lock =
new std::unique_lock<std::mutex>(*mu);
153 outF.write((
char *)_kit->
get_buffer(mem, block), block);
156 speed_count += block;
159 auto count = std::chrono::steady_clock::now() - st;
160 std::cout << std::flush <<
'\r' <<
"current write file speed: " << speed_count*1000./(count.count()) <<
"MB/s" << std::endl;
162 st = std::chrono::steady_clock::now();
166 std::cout << std::endl;
172 kit->
write(0x50010, 1);
173 kit->
write(0x50010, 0);
174 kit->
write(0x51010, 1);
175 kit->
write(0x51010, 0);
176 kit->
write(0x52010, 1);
177 kit->
write(0x52010, 0);
178 kit->
write(0x53010, 1);
179 kit->
write(0x53010, 0);
187 kit->
write(0x00004, 0);
188 kit->
write(0x00004, 1);
189 kit->
write(0x00008, 0);
190 kit->
write(0x00008, 1);
191 kit->
write(0x0000C, 0);
192 kit->
write(0x0000C, 1);
193 kit->
write(0x00010, 0);
194 kit->
write(0x00010, 1);
195 kit->
write(0x2001C, 0);
196 kit->
write(0x2001C, 1);
197 kit->
write(0x2001C, 2);
198 kit->
write(0x2001C, 3);
199 kit->
write(0x20014, 0);
200 kit->
write(0x20014, 1);
201 kit->
write(0x20010, 0);
202 kit->
write(0x20010, 1);
203 kit->
write(0x20020, 0xF);
204 kit->
write(0x20020, 0);
205 kit->
write(0x00004, 0);
206 kit->
write(0x00004, 1);
207 kit->
write(0x00008, 0);
208 kit->
write(0x00008, 1);
209 kit->
write(0x0000C, 0);
210 kit->
write(0x0000C, 1);
211 kit->
write(0x00010, 0);
212 kit->
write(0x00010, 1);
218 kit->
write(0x30000, 0x0);
219 kit->
write(0x31000, 0x0);
220 kit->
write(0x32000, 0x0);
221 kit->
write(0x33000, 0x0);
222 kit->
write(0x30008, 0x0);
223 kit->
write(0x30008, 0x1);
224 kit->
write(0x31008, 0x0);
225 kit->
write(0x31008, 0x1);
226 kit->
write(0x32008, 0x0);
227 kit->
write(0x32008, 0x1);
228 kit->
write(0x33008, 0x0);
229 kit->
write(0x33008, 0x1);
235 kit->
write(0x24 , 0x1);
236 kit->
write(0x30030, 98304);
237 kit->
write(0x3000C, 49024);
238 kit->
write(0x30004, 100);
239 kit->
write(0x30010, 0x0);
240 kit->
write(0x3002C, 0xAAAA0000);
241 kit->
write(0x3003C, 0x10);
242 kit->
write(0x30000, 0x1);
243 kit->
write(0x28 , 0x1);
244 kit->
write(0x31030, 0x4000);
245 kit->
write(0x3100C, 0x1F80);
246 kit->
write(0x31004, 0x190);
247 kit->
write(0x31010, 0x0);
248 kit->
write(0x3102C, 0xAAAA1111);
249 kit->
write(0x3103C, 0x10);
250 kit->
write(0x31000, 0x1);
251 kit->
write(0x2c , 0x1);
252 kit->
write(0x32030, 0x4000);
253 kit->
write(0x3200C, 0x1F80);
254 kit->
write(0x32004, 0x190);
255 kit->
write(0x32010, 0x0);
256 kit->
write(0x3202C, 0xAAAA2222);
257 kit->
write(0x3203C, 0x10);
258 kit->
write(0x32000, 0x1);
259 kit->
write(0x30 , 0x1);
260 kit->
write(0x33030, 0x4000);
261 kit->
write(0x3300C, 0x1F80);
262 kit->
write(0x33004, 0x190);
263 kit->
write(0x33010, 0x0);
264 kit->
write(0x3302C, 0xAAAA3333);
265 kit->
write(0x3303C, 0x10);
266 kit->
write(0x33000, 0x1);
272 for (
int i=0; i<10; i++) {
278 std::thread up_trd(
upload_thread, kit, stream_chnl, q, ds_block, total_len, mu);
281 std::cout <<
"Stream thread start successful!!!" << std::endl;
289int main(
int argc,
char *argv[]) {
290 unsigned int ds_block = 4*1024*1024;
291 Deque q0, q1, q2, q3;
292 std::mutex mu0, mu1, mu2, mu3;
296 std::cout <<
"Unsupported parameter passing method" << std::endl;
298 std::cout << argv[0] <<
" {totalBytes}" << std::endl;
301 nsuSize_t total_len = std::atoi(argv[1]);
302 if (total_len % ds_block != 0) {
303 std::cout <<
"The total length of upstream data total_len "
304 << total_len <<
" Bytes should be "
305 << ds_block <<
"Integer multiple of Bytes" << std::endl;
312#ifdef STREAM_WITH_TCP
313 auto ip = std::string(argv[1]);
315 std::string port_str{};
316 port_str += ip[ip.length()-2];
318 port_str += ip[ip.length()-1];
319 int port = std::atoi(port_str.data());
325 auto res = kit.link_cmd(¶m);
330 res = kit.link_stream(¶m);
341 std::cout <<
"SocLink Successful!!! now alloc buffer" << std::endl;
342 for (
int i=0; i<10; i++) {
347 for (
int i=0; i<10; i++) {
352 for (
int i=0; i<10; i++) {
357 for (
int i=0; i<10; i++) {
363 std::thread up_trd_0(
upload_thread, &kit, 0, &q0, ds_block, total_len, &mu0);
364 std::thread write_trd_0(
write_file_thread, &kit, &q0, ds_block,
"data_ch_0.dat", &mu0);
365 std::thread up_trd_1(
upload_thread, &kit, 1, &q1, ds_block, total_len, &mu1);
366 std::thread write_trd_1(
write_file_thread, &kit, &q1, ds_block,
"data_ch_1.dat", &mu1);
367 std::thread up_trd_2(
upload_thread, &kit, 2, &q2, ds_block, total_len, &mu2);
368 std::thread write_trd_2(
write_file_thread, &kit, &q2, ds_block,
"data_ch_2.dat", &mu2);
369 std::thread up_trd_3(
upload_thread, &kit, 3, &q3, ds_block, total_len, &mu3);
370 std::thread write_trd_3(
write_file_thread, &kit, &q3, ds_block,
"data_ch_3.dat", &mu3);
372 std::cout <<
"Stream thread start successful!!!" << std::endl;
388 std::cout <<
"Data upload completed" << std::endl;
virtual nsukitStatus_t open_recv(nsuChnlNum_t chnl, nsuMemory_p fd, nsuStreamLen_t length, nsuStreamLen_t offset=0)
virtual nsuVoidBuf_p get_buffer(nsuMemory_p fd, nsuStreamLen_t length=0)
virtual nsuMemory_p alloc_buffer(nsuStreamLen_t length, nsuVoidBuf_p buf=nullptr)
virtual nsukitStatus_t wait_stream(nsuMemory_p fd, float timeout=0.)
virtual nsukitStatus_t write(nsuRegAddr_t addr, nsuRegValue_t value)
void upload_thread(nsukit::BaseKit *_kit, int stream_chnl, Deque *q, nsuSize_t block, nsuSize_t total, std::mutex *mu)
int aurora_reset(nsukit::BaseKit *kit)
void write_file_thread(nsukit::BaseKit *_kit, Deque *q, nsuSize_t block, const std::string &path, std::mutex *mu)
int source_select(nsukit::BaseKit *kit)
int ddr_reset(nsukit::BaseKit *kit)
int start_chnl(nsukit::BaseKit *kit, Deque *q, uint32_t stream_chnl, uint32_t ds_block, uint32_t total_len, std::string file_name, std::mutex *mu)
int source_reset(nsukit::BaseKit *kit)
ThreadSafeQueue< void > memQueue
std::string NSU_DLLEXPORT status2_string(nsukitStatus_t status)
uint32_t stream_tcp_block
nsuBoardNum_t stream_board
@ NSUKIT_STATUS_STREAM_RUNNING
DLLEXTERN typedef size_t nsuSize_t
DLLEXTERN typedef void * nsuMemory_p