Commit ba98b0cf authored by Andrey Filippov's avatar Andrey Filippov
Browse files

working version with 2 parallel write threads

parent bff5f99c
Loading
Loading
Loading
Loading
+346 −125

File changed.

Preview size limit exceeded, changes collapsed.

+58 −20
Original line number Original line Diff line number Diff line
@@ -18,7 +18,7 @@


#ifndef _CAMOGM_H
#ifndef _CAMOGM_H
#define _CAMOGM_H
#define _CAMOGM_H
#define __USE_GNU // for O_DIRECT
//#define __USE_GNU // for O_DIRECT
//#define  USE_POLL
//#define  USE_POLL
#include <pthread.h>
#include <pthread.h>
#include <stdbool.h>
#include <stdbool.h>
@@ -30,6 +30,11 @@
#include <syslog.h>
#include <syslog.h>
#include <malloc.h> // debugging
#include <malloc.h> // debugging


#define NUM_WRITER_THREADS        2
#define NUM_NEED_EMPTY            (NUM_WRITER_THREADS + 1) // run main thread to prepage more pages if >= empty
#define NUM_NEED_NEMPTY           (NUM_WRITER_THREADS + 0) // wake up main thread if there are at least this number of empty pages
#define NUM_FREE_MANY             (NUM_WRITER_THREADS + 1) // Set many if get_num_empty() < NUM_FREE_MANY

//#define NODEBUG
//#define NODEBUG
#define CAMOGM_FRAME_NOT_READY    1        ///< frame pointer valid, but not yet acquired
#define CAMOGM_FRAME_NOT_READY    1        ///< frame pointer valid, but not yet acquired
#define CAMOGM_FRAME_INVALID      2        ///< invalid frame pointer
#define CAMOGM_FRAME_INVALID      2        ///< invalid frame pointer
@@ -49,6 +54,8 @@
#define CAMOGM_FORMAT_MOV         3        ///< output as Apple Quicktime
#define CAMOGM_FORMAT_MOV         3        ///< output as Apple Quicktime
#define CAMOGM_MIN_BUF_FRAMES     2        ///< buffer should accomodate at list this number of frames
#define CAMOGM_MIN_BUF_FRAMES     2        ///< buffer should accomodate at list this number of frames


#define FLUSH_DEBUG               1        ///< flush debug file after each print

#ifdef NODEBUG
#ifdef NODEBUG
#define D(x)
#define D(x)
#define D0(x)
#define D0(x)
@@ -60,7 +67,9 @@
#define D6(x)
#define D6(x)
#define D7(x)
#define D7(x)
#define DD(x)
#define DD(x)
#define DFLUSH
#else
#else
# if FLUSH_DEBUG
#define D(x) { if (debug_file && debug_level) { x; fflush(debug_file); } }
#define D(x) { if (debug_file && debug_level) { x; fflush(debug_file); } }
#define D0(x) { if (debug_file) { pthread_mutex_lock(&print_mutex); x; fflush(debug_file); pthread_mutex_unlock(&print_mutex); } }
#define D0(x) { if (debug_file) { pthread_mutex_lock(&print_mutex); x; fflush(debug_file); pthread_mutex_unlock(&print_mutex); } }
#define D1(x) { if (debug_file && (debug_level > 0)) { pthread_mutex_lock(&print_mutex); x; fflush(debug_file); pthread_mutex_unlock(&print_mutex); } }
#define D1(x) { if (debug_file && (debug_level > 0)) { pthread_mutex_lock(&print_mutex); x; fflush(debug_file); pthread_mutex_unlock(&print_mutex); } }
@@ -72,6 +81,21 @@
//#define D7(x) { if (debug_file && (debug_level > 6)) { pthread_mutex_lock(&print_mutex); x; fflush(debug_file); pthread_mutex_unlock(&print_mutex); } }
//#define D7(x) { if (debug_file && (debug_level > 6)) { pthread_mutex_lock(&print_mutex); x; fflush(debug_file); pthread_mutex_unlock(&print_mutex); } }
#define D7(x)
#define D7(x)
#define DD(x)  { if (debug_file) { fprintf(debug_file, "%s:%d:", __FILE__, __LINE__); x; fflush(debug_file); } }
#define DD(x)  { if (debug_file) { fprintf(debug_file, "%s:%d:", __FILE__, __LINE__); x; fflush(debug_file); } }
#else // if FLUSH_DEBUG
#define D(x) { if (debug_file && debug_level) { x; } }
#define D0(x) { if (debug_file)                      { pthread_mutex_lock(&print_mutex); x; pthread_mutex_unlock(&print_mutex); } }
#define D1(x) { if (debug_file && (debug_level > 0)) { pthread_mutex_lock(&print_mutex); x; pthread_mutex_unlock(&print_mutex); } }
#define D2(x) { if (debug_file && (debug_level > 1)) { pthread_mutex_lock(&print_mutex); x; pthread_mutex_unlock(&print_mutex); } }
#define D3(x) { if (debug_file && (debug_level > 2)) { pthread_mutex_lock(&print_mutex); x; pthread_mutex_unlock(&print_mutex); } }
#define D4(x) { if (debug_file && (debug_level > 3)) { pthread_mutex_lock(&print_mutex); x; pthread_mutex_unlock(&print_mutex); } }
#define D5(x) { if (debug_file && (debug_level > 4)) { pthread_mutex_lock(&print_mutex); x; pthread_mutex_unlock(&print_mutex); } }
#define D6(x) { if (debug_file && (debug_level > 5)) { pthread_mutex_lock(&print_mutex); x; pthread_mutex_unlock(&print_mutex); } }
//#define D7(x) { if (debug_file && (debug_level > 6)) { pthread_mutex_lock(&print_mutex); x; pthread_mutex_unlock(&print_mutex); } }
#define D7(x)
#define DD(x)  { if (debug_file) { fprintf(debug_file, "%s:%d:", __FILE__, __LINE__); x; } }
#endif // if FLUSH_DEBUG else
#define DFLUSH { if (debug_file ) { pthread_mutex_lock(&print_mutex);fflush(debug_file); pthread_mutex_unlock(&print_mutex); } }

#endif
#endif


/** @brief HEADER_SIZE is defined to be larger than actual header (with EXIF) to use compile-time buffer */
/** @brief HEADER_SIZE is defined to be larger than actual header (with EXIF) to use compile-time buffer */
@@ -92,6 +116,8 @@
#define JPEG_TRAILER_LEN          2  ///< The size in bytes of JPEG trailer
#define JPEG_TRAILER_LEN          2  ///< The size in bytes of JPEG trailer
#define CIRCBUF_ALIGNMENT_SIZE   32  ///< Align of CIRCBUF entries
#define CIRCBUF_ALIGNMENT_SIZE   32  ///< Align of CIRCBUF entries


#define MIN_USED_SIZE      4000000 // debugging DMA clashes

// Switching CHUNK -> SEGMENT
// Switching CHUNK -> SEGMENT
#define PAGE_PHYS              4096 //  512 // alignment size for O_DIRECT, may be different from LBA block size that is 512
#define PAGE_PHYS              4096 //  512 // alignment size for O_DIRECT, may be different from LBA block size that is 512
enum segments  {
enum segments  {
@@ -101,7 +127,13 @@ enum segments {
	SEGMENTS_NUMBER     ///< Just to get number of segment types
	SEGMENTS_NUMBER     ///< Just to get number of segment types
};
};


#define SEGMENTS_PAGES         4 // 2
enum segpage_states {
	SEGPAGE_EMPTY = 0,  ///< page of segments is empty (new or written to disk)
	SEGPAGE_FULL,       ///< page of segments is prepared for writing
	SEGPAGE_BUSY        ///< page of segments is in the process of being sent to fisk
};

#define SEGMENTS_PAGES         8 // 4 // 2
#define SEGMENTS_TOTAL         ((SEGMENTS_NUMBER) * (SEGMENTS_PAGES))
#define SEGMENTS_TOTAL         ((SEGMENTS_NUMBER) * (SEGMENTS_PAGES))


/** Glue buffer should be large enough to contain:
/** Glue buffer should be large enough to contain:
@@ -203,29 +235,35 @@ typedef struct {
 * @brief Contains mutexes and conditional variables associated with disk writing thread
 * @brief Contains mutexes and conditional variables associated with disk writing thread
 */
 */
struct writer_params {
struct writer_params {
	int blockdev_fd;                                        ///< file descriptor for open block device where frame will be recorded
//	int blockdev_fd;                                        ///< file descriptor for open block device where frame will be recorded
	pthread_t writer_thread;                                ///< disk writing thread
	pthread_t writer_threads [NUM_WRITER_THREADS];          ///< array of disk writer threads
	bool writev_run          [NUM_WRITER_THREADS];          ///< writev() is active (enable main thread if buffer not full and writew or buf empty)
	bool write_waits_sig     [NUM_WRITER_THREADS];          ///< This write thread is waiting for signal
	bool write_go            [NUM_WRITER_THREADS];          ///< This writer thread may go even if other thread already left write_waits_sig state
	int writev_run_segm      [NUM_WRITER_THREADS];          ///< current writev segment being written per thread (normally segment0 is small, other - big

	pthread_mutex_t writer_mutex;                           ///< synchronization mutex for main and writing threads
	pthread_mutex_t writer_mutex;                           ///< synchronization mutex for main and writing threads
	pthread_cond_t writer_cond;                             ///< conditional variable indicating that writer thread can proceed with new frame
	pthread_cond_t writer_cond;                             ///< conditional variable indicating that writer thread can proceed with new frame
	pthread_cond_t main_cond;                               ///< conditional variable indicating that main thread can update write pointers
	pthread_cond_t main_cond;                               ///< conditional variable indicating that main thread can update write pointers
	bool circbuf_many;                                      ///< probably there are many more frames in circbuf
    int chunk_page_prep;                                    ///< page of chunks being prepared. Incremented (mod) after data is prepared for raw write
    int chunk_page_prep;                                    ///< page of chunks being prepared. Incremented (mod) after data is prepared for raw write
    int chunk_page_write;                                   ///< page of chunks to write. Incremented (mod) after recording to disk
    int chunk_page_state[SEGMENTS_PAGES];                   ///< segments page state: 0 - empty, 1 - full, 2 - currently writing to disk
	bool writev_run;                                        ///< writev() is active (enable main thread if buffer not full and writew or buf empty)
    uint64_t next_segment_pos;                              ///< next segment calculated file position (bytes)
    uint64_t segment_pos [SEGMENTS_PAGES];                  ///< per-segment file offsets
	int last_ret_val;                                       ///< error value return during last frame recording (if any occurred)
	int last_ret_val;                                       ///< error value return during last frame recording (if any occurred)
	bool exit_thread;                                       ///< flag indicating that the writing thread should terminate
	bool exit_write_threads;                                ///< flag indicating that the writing threads should terminate and close their files
	int state;                                              ///< the state of disk writing thread
	int state;                                              ///< the state of disk writing thread
	struct iovec *data_segments[SEGMENTS_PAGES];            ///< a set of vectors sets pointing to aligned frame data buffers
	struct iovec *data_segments[SEGMENTS_PAGES];            ///< a set of vectors sets pointing to aligned frame data buffers
	unsigned char *common_buffs[FILE_CHUNKS_PAGES];         ///< buffer for aligned JPEG header // make multiple?
	int  dbg_data [SEGMENTS_PAGES][2];                      ///< debug data: channel and length of copied tail
	unsigned char *glue_buffs[SEGMENTS_PAGES + 1];          ///< glue buffer (end of previous frame + start of this one
	unsigned char *glue_buffs[SEGMENTS_PAGES + 1];          ///< glue buffer (end of previous frame + start of this one
//	unsigned char *glue_buffs_carry;                        ///< extra glue buffer to pass from the previous to the next frame
	struct iovec  glue_carry_vec;                           ///< current tail pointer of the carry glue segment
	struct iovec  glue_carry_vec;                           ///< current tail pointer of the carry glue segment
	uint64_t lba_start;                                     ///< disk starting LBA
	uint64_t lba_start;                                     ///< disk starting LBA
	uint64_t lba_current;                                   ///< current write position in LBAs
	uint64_t lba_end;                                       ///< disk last LBA
	uint64_t lba_end;                                       ///< disk last LBA

	uint64_t stat_update;                                   ///< time when status file was updated
	time_t stat_update;                                     ///< time when status file was updated
	bool dummy_read;                                        ///< enable dummy read cycle (debug feature)
	bool dummy_read;                                        ///< enable dummy read cycle (debug feature)
};
};


/**
/**
 * @struct camogm_state
 * @struct camogm_state
 * @brief Holds current state of the running program
 * @brief Holds current state of the running program
@@ -337,8 +375,8 @@ int waitDaemonEnabled(unsigned int port, int daemonBit);
int isDaemonEnabled(unsigned int port, int daemonBit);
int isDaemonEnabled(unsigned int port, int daemonBit);
int is_fd_valid(int fd);
int is_fd_valid(int fd);
int get_fpga_usec(const int fd_fparsall, unsigned int port);
int get_fpga_usec(const int fd_fparsall, unsigned int port);
inline void wait_frame_sync(const int fd_fparsall);
uint64_t get_fpga_time64(const int fd_fparsall, unsigned int port);
void wait_frame_sync(const int fd_fparsall);
unsigned long *get_ccam_dma_buf(int port);
unsigned long *get_ccam_dma_buf(int port);



#endif /* _CAMOGM_H */
#endif /* _CAMOGM_H */
+79 −59
Original line number Original line Diff line number Diff line
@@ -95,7 +95,7 @@ static inline unsigned char *vectrpos(struct iovec *vec, size_t offset)
/** extend to page-aligned */
/** extend to page-aligned */
int vectaligntail(struct iovec *dest)
int vectaligntail(struct iovec *dest)
{
{
	unsigned char *d = (unsigned char *)dest->iov_base;
//	unsigned char *d = (unsigned char *)dest->iov_base;
	int len = dest->iov_len;
	int len = dest->iov_len;
	int ceil_size = (len & (PAGE_PHYS - 1)) ? ((len | (PAGE_PHYS - 1)) + 1) : len;
	int ceil_size = (len & (PAGE_PHYS - 1)) ? ((len | (PAGE_PHYS - 1)) + 1) : len;
	if (ceil_size > len) {
	if (ceil_size > len) {
@@ -116,22 +116,23 @@ size_t remap_vectors(camogm_state *state) // , struct iovec *chunks)
	// state->writer_params.glue_carry_vec has data for the the trailer of the previous frame. swap it with the
	// state->writer_params.glue_carry_vec has data for the the trailer of the previous frame. swap it with the
	// state->writer_params.glue_buffs[seg_page][SEGMENT_GLUE]
	// state->writer_params.glue_buffs[seg_page][SEGMENT_GLUE]
	size_t total_sz;
	size_t total_sz;
	int seg_page = (state->writer_params.chunk_page_prep % FILE_CHUNKS_PAGES);
	struct writer_params *params = &state->writer_params;
	int seg_page = params->chunk_page_prep;
	long head_bytes, tail_bytes, trimmed_bytes;
	long head_bytes, tail_bytes, trimmed_bytes;
	int  i, data_index, data_pre, num_glue_pages, trim_segment, circbuff_offs;
	int  i, data_index, data_pre, num_glue_pages, trim_segment, circbuff_offs;
	int gap_size;
	int gap_size;
	unsigned char * aligned_data;
	unsigned char * aligned_data;
	struct iovec  glue_vec;
	struct iovec  glue_vec;
	struct iovec  * glue_vec_p; //  = &glue_vec;
	struct iovec  * glue_vec_p; //  = &glue_vec;
	glue_vec.iov_base = state->writer_params.glue_carry_vec.iov_base;
	glue_vec.iov_base = params->glue_carry_vec.iov_base;
	glue_vec.iov_len = state->writer_params.glue_carry_vec.iov_len;
	glue_vec.iov_len = params->glue_carry_vec.iov_len;
	D7(fprintf(debug_file, "_13c01_:remap_vectors @ %07d\n",get_fpga_usec(state->fd_fparmsall[0], 0)));
	D7(fprintf(debug_file, "_13c01_:remap_vectors @ %07d\n",get_fpga_usec(state->fd_fparmsall[0], 0)));
	state->writer_params.glue_carry_vec.iov_base = state->writer_params.data_segments[seg_page][SEGMENT_GLUE].iov_base;
	params->glue_carry_vec.iov_base = params->data_segments[seg_page][SEGMENT_GLUE].iov_base;
	state->writer_params.glue_carry_vec.iov_len = 0; // state->writer_params.data_segments[seg_page][SEGMENT_GLUE].iov_len;
	params->glue_carry_vec.iov_len = 0; // params->data_segments[seg_page][SEGMENT_GLUE].iov_len;
	state->writer_params.data_segments[seg_page][SEGMENT_GLUE].iov_base =  glue_vec.iov_base;
	params->data_segments[seg_page][SEGMENT_GLUE].iov_base =  glue_vec.iov_base;
	state->writer_params.data_segments[seg_page][SEGMENT_GLUE].iov_len =   glue_vec.iov_len;
	params->data_segments[seg_page][SEGMENT_GLUE].iov_len =   glue_vec.iov_len;
	state->writer_params.data_segments[seg_page][SEGMENT_FIRST].iov_len =  0; // check iov_len first before using iov_base
	params->data_segments[seg_page][SEGMENT_FIRST].iov_len =  0; // check iov_len first before using iov_base
	state->writer_params.data_segments[seg_page][SEGMENT_SECOND].iov_len = 0; // check iov_len first before using iov_base
	params->data_segments[seg_page][SEGMENT_SECOND].iov_len = 0; // check iov_len first before using iov_base
	// calculate total header size
	// calculate total header size
	head_bytes = 0;
	head_bytes = 0;
	for (i = 0; i < state->chunk_data_index; i++){
	for (i = 0; i < state->chunk_data_index; i++){
@@ -146,14 +147,14 @@ size_t remap_vectors(camogm_state *state) // , struct iovec *chunks)
	circbuff_offs = (int) state->packetchunks[data_index].chunk;
	circbuff_offs = (int) state->packetchunks[data_index].chunk;
	aligned_data = state->packetchunks[data_index].chunk+ (PAGE_PHYS_CEIL(circbuff_offs) - circbuff_offs);
	aligned_data = state->packetchunks[data_index].chunk+ (PAGE_PHYS_CEIL(circbuff_offs) - circbuff_offs);
	data_pre = aligned_data - state->packetchunks[data_index].chunk;
	data_pre = aligned_data - state->packetchunks[data_index].chunk;
	glue_vec_p = &state->writer_params.data_segments[seg_page][SEGMENT_GLUE];
	glue_vec_p = &params->data_segments[seg_page][SEGMENT_GLUE];
	D7(fprintf(debug_file, "_13c03_:remap_vectors, circbuff_offs=0x%x aligned_data=%p data_pre=0x%x @ %07d\n", circbuff_offs, aligned_data, data_pre, get_fpga_usec(state->fd_fparmsall[0], 0)));
	D7(fprintf(debug_file, "_13c03_:remap_vectors, circbuff_offs=0x%x aligned_data=%p data_pre=0x%x @ %07d\n", circbuff_offs, aligned_data, data_pre, get_fpga_usec(state->fd_fparmsall[0], 0)));
	D7(fprintf(debug_file, "_13c04_:remap_vectors, state->packetchunks[state->chunk_data_index].bytes=%ld @ %07d\n", state->packetchunks[state->chunk_data_index].bytes, get_fpga_usec(state->fd_fparmsall[0], 0)));
	D7(fprintf(debug_file, "_13c04_:remap_vectors, state->packetchunks[state->chunk_data_index].bytes=%ld @ %07d\n", state->packetchunks[state->chunk_data_index].bytes, get_fpga_usec(state->fd_fparmsall[0], 0)));
	if (state->packetchunks[state->chunk_data_index].bytes < data_pre) {
	if (state->packetchunks[state->chunk_data_index].bytes < data_pre) {
		// Too short frame data - special treatment with zero data, moving some to carry, and possibly outputing some from the GLUE
		// Too short frame data - special treatment with zero data, moving some to carry, and possibly outputing some from the GLUE
		// there is definitely no second data - otherwise the first segment would end at page boundary
		// there is definitely no second data - otherwise the first segment would end at page boundary
		// copy everything to state->writer_params.data_segments[seg_page][SEGMENT_GLUE], trim to pages, pass residual to
		// copy everything to params->data_segments[seg_page][SEGMENT_GLUE], trim to pages, pass residual to
		// state->writer_params.glue_carry_vec
		// params->glue_carry_vec
		// copy header data, no gap
		// copy header data, no gap
		for (i = 0; i <= state->chunk_index; i++){ // including data and trailer
		for (i = 0; i <= state->chunk_index; i++){ // including data and trailer
			vectcpy(glue_vec_p, state->packetchunks[i].chunk, state->packetchunks[i].bytes);
			vectcpy(glue_vec_p, state->packetchunks[i].chunk, state->packetchunks[i].bytes);
@@ -161,7 +162,7 @@ size_t remap_vectors(camogm_state *state) // , struct iovec *chunks)
		trimmed_bytes = PAGE_PHYS_TRIM(glue_vec_p->iov_len);
		trimmed_bytes = PAGE_PHYS_TRIM(glue_vec_p->iov_len);
		D7(fprintf(debug_file, "_13c05_:remap_vectors, trimmed_bytes=%ld @ %07d\n", trimmed_bytes, get_fpga_usec(state->fd_fparmsall[0], 0)));
		D7(fprintf(debug_file, "_13c05_:remap_vectors, trimmed_bytes=%ld @ %07d\n", trimmed_bytes, get_fpga_usec(state->fd_fparmsall[0], 0)));
		if ((glue_vec_p->iov_len - trimmed_bytes) > 0){
		if ((glue_vec_p->iov_len - trimmed_bytes) > 0){
			vectcpy(&state->writer_params.glue_carry_vec,
			vectcpy(&params->glue_carry_vec,
					((char *) glue_vec_p->iov_base) + trimmed_bytes,
					((char *) glue_vec_p->iov_base) + trimmed_bytes,
					glue_vec_p->iov_len - trimmed_bytes);
					glue_vec_p->iov_len - trimmed_bytes);
			glue_vec_p->iov_len = trimmed_bytes; // may be 0
			glue_vec_p->iov_len = trimmed_bytes; // may be 0
@@ -178,7 +179,7 @@ size_t remap_vectors(camogm_state *state) // , struct iovec *chunks)
		D7(fprintf(debug_file, "_13c06a0_:malloc_usable_size(glue_vec_p->iov_base) = %d @ %07d\n", malloc_usable_size(glue_vec_p->iov_base), get_fpga_usec(state->fd_fparmsall[0], 0)));
		D7(fprintf(debug_file, "_13c06a0_:malloc_usable_size(glue_vec_p->iov_base) = %d @ %07d\n", malloc_usable_size(glue_vec_p->iov_base), get_fpga_usec(state->fd_fparmsall[0], 0)));
		for (i = 0; i <  SEGMENTS_PAGES; i++){
		for (i = 0; i <  SEGMENTS_PAGES; i++){
			D7(fprintf(debug_file, "_13c06a1_: malloc_usable_size([%d][SEGMENT_GLUE]) = %d @ %07d\n",\
			D7(fprintf(debug_file, "_13c06a1_: malloc_usable_size([%d][SEGMENT_GLUE]) = %d @ %07d\n",\
					i, malloc_usable_size(state->writer_params.data_segments[i][SEGMENT_GLUE].iov_base), get_fpga_usec(state->fd_fparmsall[0], 0)));
					i, malloc_usable_size(params->data_segments[i][SEGMENT_GLUE].iov_base), get_fpga_usec(state->fd_fparmsall[0], 0)));
		}
		}


		// size_t malloc_usable_size(void *ptr);
		// size_t malloc_usable_size(void *ptr);
@@ -197,44 +198,49 @@ size_t remap_vectors(camogm_state *state) // , struct iovec *chunks)
		vectcpy(glue_vec_p, state->packetchunks[state->chunk_data_index].chunk, data_pre);
		vectcpy(glue_vec_p, state->packetchunks[state->chunk_data_index].chunk, data_pre);
		// TODO: verify that glue_vec_p->iov_len is multiple of PAGE_PHYS
		// TODO: verify that glue_vec_p->iov_len is multiple of PAGE_PHYS
		trim_segment = SEGMENT_FIRST;
		trim_segment = SEGMENT_FIRST;
		state->writer_params.data_segments[seg_page][SEGMENT_FIRST].iov_base = ((char *) state->packetchunks[state->chunk_data_index].chunk) + data_pre;
		params->data_segments[seg_page][SEGMENT_FIRST].iov_base = ((char *) state->packetchunks[state->chunk_data_index].chunk) + data_pre;
		state->writer_params.data_segments[seg_page][SEGMENT_FIRST].iov_len =            state->packetchunks[state->chunk_data_index].bytes - data_pre;
		params->data_segments[seg_page][SEGMENT_FIRST].iov_len =            state->packetchunks[state->chunk_data_index].bytes - data_pre;
		D7(fprintf(debug_file, "_13c07_:remap_vectors @ %07d\n", get_fpga_usec(state->fd_fparmsall[0], 0)));
		D7(fprintf(debug_file, "_13c07_:remap_vectors @ %07d\n", get_fpga_usec(state->fd_fparmsall[0], 0)));
		if (state->writer_params.data_segments[seg_page][SEGMENT_FIRST].iov_len == 0) {
		if (params->data_segments[seg_page][SEGMENT_FIRST].iov_len == 0) {
			// nothing left of data segment 1 - move segment2 to segment1 (or skip both)
			// nothing left of data segment 1 - move segment2 to segment1 (or skip both)
			if (state->data_segments > 1) { // move second segment to first
			if (state->data_segments > 1) { // move second segment to first
				state->writer_params.data_segments[seg_page][SEGMENT_FIRST].iov_base = ((char *) state->packetchunks[state->chunk_data_index + 1].chunk);
				params->data_segments[seg_page][SEGMENT_FIRST].iov_base = ((char *) state->packetchunks[state->chunk_data_index + 1].chunk);
				state->writer_params.data_segments[seg_page][SEGMENT_FIRST].iov_len =            state->packetchunks[state->chunk_data_index + 1].bytes ;
				params->data_segments[seg_page][SEGMENT_FIRST].iov_len =            state->packetchunks[state->chunk_data_index + 1].bytes ;
			}
			}
		} else if (state->data_segments > 1) {
		} else if (state->data_segments > 1) {
			state->writer_params.data_segments[seg_page][SEGMENT_SECOND].iov_base = ((char *) state->packetchunks[state->chunk_data_index + 1].chunk);
			params->data_segments[seg_page][SEGMENT_SECOND].iov_base = ((char *) state->packetchunks[state->chunk_data_index + 1].chunk);
			state->writer_params.data_segments[seg_page][SEGMENT_SECOND].iov_len =            state->packetchunks[state->chunk_data_index + 1].bytes ;
			params->data_segments[seg_page][SEGMENT_SECOND].iov_len =            state->packetchunks[state->chunk_data_index + 1].bytes ;
			trim_segment = SEGMENT_SECOND;
			trim_segment = SEGMENT_SECOND;
		}
		}
		D7(fprintf(debug_file, "_13c08_:remap_vectors,[SEGMENT_GLUE].iov_len=%d, [SEGMENT_FIRST].iov_len=%d, [SEGMENT_SECOND].iov_len=%d  @ %07d\n", \
		D7(fprintf(debug_file, "_13c08_:remap_vectors,[SEGMENT_GLUE].iov_len=%d, [SEGMENT_FIRST].iov_len=%d, [SEGMENT_SECOND].iov_len=%d  @ %07d\n", \
				state->writer_params.data_segments[seg_page][SEGMENT_GLUE].iov_len, \
				params->data_segments[seg_page][SEGMENT_GLUE].iov_len, \
				state->writer_params.data_segments[seg_page][SEGMENT_FIRST].iov_len, \
				params->data_segments[seg_page][SEGMENT_FIRST].iov_len, \
				state->writer_params.data_segments[seg_page][SEGMENT_SECOND].iov_len, \
				params->data_segments[seg_page][SEGMENT_SECOND].iov_len, \
				get_fpga_usec(state->fd_fparmsall[0], 0)));
				get_fpga_usec(state->fd_fparmsall[0], 0)));


		// Trim last data segment to PAGE_PHYS (it may make its length zero - OK), copy remaining
		// Trim last data segment to PAGE_PHYS (it may make its length zero - OK), copy remaining
		// data to the buffer pointed by state->writer_params.glue_carry_vec.iov_base, set state->writer_params.glue_carry_vec.iov_len
		// data to the buffer pointed by params->glue_carry_vec.iov_base, set params->glue_carry_vec.iov_len
		trimmed_bytes = PAGE_PHYS_TRIM(state->writer_params.data_segments[seg_page][trim_segment].iov_len);
		trimmed_bytes = PAGE_PHYS_TRIM(params->data_segments[seg_page][trim_segment].iov_len);
		D7(fprintf(debug_file, "_13c09_:remap_vectors, trimmed_bytes=%ld @ %07d\n", trimmed_bytes, get_fpga_usec(state->fd_fparmsall[0], 0)));
		D7(fprintf(debug_file, "_13c09_:remap_vectors, trimmed_bytes=%ld @ %07d\n", trimmed_bytes, get_fpga_usec(state->fd_fparmsall[0], 0)));
		if (state->writer_params.data_segments[seg_page][trim_segment].iov_len > trimmed_bytes) {
		// debugging dma-related bugs.
		params->dbg_data[seg_page][0] = state->port_num;
		params->dbg_data[seg_page][1] = params->data_segments[seg_page][trim_segment].iov_len - trimmed_bytes; // remainder copied to glue buffer
		if (params->data_segments[seg_page][trim_segment].iov_len > trimmed_bytes) {
			D6(fprintf(debug_file, "_13c09a:remap_vectors, trim_segment=%d, dest = %p, dest_len=%d src=%p, copy_len = %ld @ %07d\n", \
			D6(fprintf(debug_file, "_13c09a:remap_vectors, trim_segment=%d, dest = %p, dest_len=%d src=%p, copy_len = %ld @ %07d\n", \
					trim_segment,
					trim_segment,
					&state->writer_params.glue_carry_vec, \
					&params->glue_carry_vec, \
					state->writer_params.glue_carry_vec.iov_len, \
					params->glue_carry_vec.iov_len, \
					((char *) state->writer_params.data_segments[seg_page][trim_segment].iov_base) + trimmed_bytes, \
					((char *) params->data_segments[seg_page][trim_segment].iov_base) + trimmed_bytes, \
					state->writer_params.data_segments[seg_page][trim_segment].iov_len - trimmed_bytes, \
					params->data_segments[seg_page][trim_segment].iov_len - trimmed_bytes, \
					get_fpga_usec(state->fd_fparmsall[0], 0)));
					get_fpga_usec(state->fd_fparmsall[0], 0)));


			vectcpy(&state->writer_params.glue_carry_vec,
			vectcpy(&params->glue_carry_vec,
					((char *) state->writer_params.data_segments[seg_page][trim_segment].iov_base) + trimmed_bytes,
					((char *) params->data_segments[seg_page][trim_segment].iov_base) + trimmed_bytes,
					state->writer_params.data_segments[seg_page][trim_segment].iov_len - trimmed_bytes);
					params->data_segments[seg_page][trim_segment].iov_len - trimmed_bytes);
			state->writer_params.data_segments[seg_page][trim_segment].iov_len = trimmed_bytes; // reduce data segment length
			params->data_segments[seg_page][trim_segment].iov_len = trimmed_bytes; // reduce data segment length
		}
		}


		D7(fprintf(debug_file, "_13c09b_:remap_vectors, tail_segm=%d, chunk_index=%d @ %07d\n", \
		D7(fprintf(debug_file, "_13c09b_:remap_vectors, tail_segm=%d, chunk_index=%d @ %07d\n", \
				state->chunk_data_index + state->data_segments, \
				state->chunk_data_index + state->data_segments, \
				state->chunk_index, \
				state->chunk_index, \
@@ -243,13 +249,13 @@ size_t remap_vectors(camogm_state *state) // , struct iovec *chunks)
		for (i = state->chunk_data_index + state->data_segments; i < state->chunk_index; i++){ // trailer if it exists
		for (i = state->chunk_data_index + state->data_segments; i < state->chunk_index; i++){ // trailer if it exists
			D7(fprintf(debug_file, "_13c09c_:remap_vectors, i=%d, dest = %p, dest_len=%d src=%p, len =%ld @ %07d\n", \
			D7(fprintf(debug_file, "_13c09c_:remap_vectors, i=%d, dest = %p, dest_len=%d src=%p, len =%ld @ %07d\n", \
					i, \
					i, \
					&state->writer_params.glue_carry_vec, \
					&params->glue_carry_vec, \
					state->writer_params.glue_carry_vec.iov_len, \
					params->glue_carry_vec.iov_len, \
					state->packetchunks[i].chunk, \
					state->packetchunks[i].chunk, \
					state->packetchunks[i].bytes, \
					state->packetchunks[i].bytes, \
					get_fpga_usec(state->fd_fparmsall[0], 0)));
					get_fpga_usec(state->fd_fparmsall[0], 0)));


			vectcpy(&state->writer_params.glue_carry_vec,
			vectcpy(&params->glue_carry_vec,
					state->packetchunks[i].chunk,
					state->packetchunks[i].chunk,
					state->packetchunks[i].bytes);
					state->packetchunks[i].bytes);
		}
		}
@@ -257,7 +263,7 @@ size_t remap_vectors(camogm_state *state) // , struct iovec *chunks)
	D7(fprintf(debug_file, "_13c10_:remap_vectors @ %07d\n", get_fpga_usec(state->fd_fparmsall[0], 0)));
	D7(fprintf(debug_file, "_13c10_:remap_vectors @ %07d\n", get_fpga_usec(state->fd_fparmsall[0], 0)));
	total_sz = 0; // get_size_from_paged(state, 0, INCLUDE_REM) + rbuff->iov_len;
	total_sz = 0; // get_size_from_paged(state, 0, INCLUDE_REM) + rbuff->iov_len;
	for (i = 0; i < SEGMENTS_NUMBER; i++){
	for (i = 0; i < SEGMENTS_NUMBER; i++){
		total_sz += state->writer_params.data_segments[seg_page][i].iov_len;
		total_sz += params->data_segments[seg_page][i].iov_len;
	}
	}
	//port_num
	//port_num
	state->image_size[state->port_num] = total_sz;
	state->image_size[state->port_num] = total_sz;
@@ -268,30 +274,31 @@ size_t remap_vectors(camogm_state *state) // , struct iovec *chunks)
/** Allocate and initialize buffers for frame alignment */
/** Allocate and initialize buffers for frame alignment */
int init_align_buffers(camogm_state *state)
int init_align_buffers(camogm_state *state)
{
{
    D3(fprintf(debug_file, "GLUE_SEGMENTS = %d pages of %d bytes each\n", GLUE_SEGMENTS, PAGE_PHYS));
	struct writer_params *params = &state->writer_params;
	int seg_page;
	int seg_page;
    D3(fprintf(debug_file, "GLUE_SEGMENTS = %d pages of %d bytes each\n", GLUE_SEGMENTS, PAGE_PHYS));
	for (seg_page = 0; seg_page < SEGMENTS_PAGES + 1; seg_page++){
	for (seg_page = 0; seg_page < SEGMENTS_PAGES + 1; seg_page++){
		// allocate SEGMENT_GLUE buffer that will contain data copied from non-pagfe aligned segments
		// allocate SEGMENT_GLUE buffer that will contain data copied from non-pagfe aligned segments
		state->writer_params.glue_buffs[seg_page] = (unsigned char *) aligned_alloc (PAGE_PHYS, GLUE_SEGMENTS * PAGE_PHYS);
		params->glue_buffs[seg_page] = (unsigned char *) aligned_alloc (PAGE_PHYS, GLUE_SEGMENTS * PAGE_PHYS);
		if (state->writer_params.glue_buffs[seg_page] == NULL) {
		if (params->glue_buffs[seg_page] == NULL) {
			return -1;
			return -1;
		}
		}
		D6(fprintf(debug_file, "init_align_buffers(), state->writer_params.glue_buffs[%d]=%p @ %07d\n", \
		D6(fprintf(debug_file, "init_align_buffers(), params->glue_buffs[%d]=%p @ %07d\n", \
				seg_page, state->writer_params.glue_buffs[seg_page], get_fpga_usec(state->fd_fparmsall[0], 0)));
				seg_page, params->glue_buffs[seg_page], get_fpga_usec(state->fd_fparmsall[0], 0)));
	}
	}
	for (seg_page = 0; seg_page < SEGMENTS_PAGES; seg_page++){
	for (seg_page = 0; seg_page < SEGMENTS_PAGES; seg_page++){
		state->writer_params.data_segments[seg_page] = (struct iovec *)malloc(SEGMENTS_NUMBER * sizeof(struct iovec));
		params->data_segments[seg_page] = (struct iovec *)malloc(SEGMENTS_NUMBER * sizeof(struct iovec));
		if (state->writer_params.data_segments[seg_page] == NULL) {
		if (params->data_segments[seg_page] == NULL) {
			return -1;
			return -1;
		}
		}
		D6(fprintf(debug_file, "init_align_buffers(), seg_page=%d, state->writer_params.data_segments[seg_page]=%p @ %07d\n", \
		D6(fprintf(debug_file, "init_align_buffers(), seg_page=%d, params->data_segments[seg_page]=%p @ %07d\n", \
				seg_page, state->writer_params.data_segments[seg_page], get_fpga_usec(state->fd_fparmsall[0], 0)));
				seg_page, params->data_segments[seg_page], get_fpga_usec(state->fd_fparmsall[0], 0)));


		state->writer_params.data_segments[seg_page][SEGMENT_GLUE].iov_base = state->writer_params.glue_buffs[seg_page];
		params->data_segments[seg_page][SEGMENT_GLUE].iov_base = params->glue_buffs[seg_page];
		state->writer_params.data_segments[seg_page][SEGMENT_GLUE].iov_len =  0;
		params->data_segments[seg_page][SEGMENT_GLUE].iov_len =  0;
	}
	}
	state->writer_params.glue_carry_vec.iov_base = state->writer_params.glue_buffs[SEGMENTS_PAGES];
	params->glue_carry_vec.iov_base = params->glue_buffs[SEGMENTS_PAGES];
	state->writer_params.glue_carry_vec.iov_len = 0; // nothing there fro the first frame (TODO: need same for start write?)
	params->glue_carry_vec.iov_len = 0; // nothing there fro the first frame (TODO: need same for start write?)


	return 0;
	return 0;
}
}
@@ -300,25 +307,27 @@ int init_align_buffers(camogm_state *state)
void deinit_align_buffers(camogm_state *state)
void deinit_align_buffers(camogm_state *state)
{
{
	int seg_page;
	int seg_page;
	struct writer_params *params = &state->writer_params;
	for (seg_page = 0; seg_page < SEGMENTS_PAGES+1; seg_page++){
	for (seg_page = 0; seg_page < SEGMENTS_PAGES+1; seg_page++){
		free (state->writer_params.glue_buffs[seg_page]);
		free (params->glue_buffs[seg_page]);
		state->writer_params.glue_buffs[seg_page] = NULL;
		params->glue_buffs[seg_page] = NULL;
	}
	}
	for (seg_page = 0; seg_page < SEGMENTS_PAGES; seg_page++){
	for (seg_page = 0; seg_page < SEGMENTS_PAGES; seg_page++){
		free(state->writer_params.data_segments[seg_page]);
		free(params->data_segments[seg_page]);
		state->writer_params.data_segments[seg_page] = NULL;
		params->data_segments[seg_page] = NULL;
	}
	}
}
}


void reset_segments(camogm_state *state, int all, int page)
void reset_segments(camogm_state *state, int all, int page)
{
{
	int i;
	int i;
	struct iovec *vects = state->writer_params.data_segments[page];
	struct writer_params *params = &state->writer_params;
	struct iovec *vects = params->data_segments[page];
	for (i = 0; i < SEGMENTS_NUMBER; i++) {
	for (i = 0; i < SEGMENTS_NUMBER; i++) {
		vects[i].iov_len = 0;
		vects[i].iov_len = 0;
	}
	}
	if (all) {
	if (all) {
		state->writer_params.glue_carry_vec.iov_len = 0;
		params->glue_carry_vec.iov_len = 0;
	}
	}
}
}


@@ -328,3 +337,14 @@ uint64_t lba_to_offset(uint64_t lba)
{
{
	return lba * PHY_BLOCK_SIZE;
	return lba * PHY_BLOCK_SIZE;
}
}

/** Get next LBA after segment prepared for writing. There is no actual
 * write pointer as recording can be out of order. When the "file" is closed,
 * this pointer matches the normal file pointer, just in blocks.
 * @param params record state parameters
 * @return full LBA (from the disk start) write pointer, measured in 512-byte sectors.
 */
uint64_t get_lba_next(const struct writer_params *params){
	return params->lba_start + (params->next_segment_pos / PHY_BLOCK_SIZE); // next prepared (not yet written)
}
+1 −0
Original line number Original line Diff line number Diff line
@@ -49,6 +49,7 @@ void deinit_align_buffers(camogm_state *state);
void reset_segments(camogm_state *state, int all, int page);
void reset_segments(camogm_state *state, int all, int page);
size_t remap_vectors(camogm_state *state); //, struct iovec *chunks);
size_t remap_vectors(camogm_state *state); //, struct iovec *chunks);
uint64_t lba_to_offset(uint64_t lba);
uint64_t lba_to_offset(uint64_t lba);
uint64_t get_lba_next(const struct writer_params *params);
int vectaligntail(struct iovec *dest);
int vectaligntail(struct iovec *dest);




+570 −177

File changed.

Preview size limit exceeded, changes collapsed.

Loading