diff --git a/sesman/chansrv/audin.c b/sesman/chansrv/audin.c index 9c59310c..a8900a73 100644 --- a/sesman/chansrv/audin.c +++ b/sesman/chansrv/audin.c @@ -77,7 +77,6 @@ static struct xr_wave_format_ex g_pcm_44100 = static struct chansrv_drdynvc_procs g_audin_info; static int g_audin_chanid; -static struct stream *g_in_s; static struct xr_wave_format_ex *g_server_formats[] = { @@ -422,64 +421,9 @@ audin_close_response(int chan_id) LOG_DEVEL(LOG_LEVEL_INFO, "audin_close_response:"); g_audin_chanid = 0; cleanup_client_formats(); - free_stream(g_in_s); - g_in_s = NULL; return 0; } -/*****************************************************************************/ -static int -audin_data_fragment(int chan_id, struct stream *s) -{ - int rv; - int bytes = s_rem(s); - LOG_DEVEL(LOG_LEVEL_DEBUG, "audin_data_fragment:"); - if (!s_check_rem(g_in_s, bytes)) - { - LOG_DEVEL(LOG_LEVEL_ERROR, "audin_data_fragment: error bytes %d left %d", - bytes, (int) (g_in_s->end - g_in_s->p)); - return 1; - } - out_uint8a(g_in_s, s->p, bytes); - if (g_in_s->p == g_in_s->end) - { - g_in_s->p = g_in_s->data; - rv = audin_process_msg(chan_id, g_in_s); - free_stream(g_in_s); - g_in_s = NULL; - return rv; - } - return 0; -} - -/*****************************************************************************/ -static int -audin_data_first(int chan_id, struct stream *s, int total_bytes) -{ - LOG_DEVEL(LOG_LEVEL_DEBUG, "audin_data_first:"); - if (g_in_s != NULL) - { - LOG_DEVEL(LOG_LEVEL_ERROR, "audin_data_first: warning g_in_s is not nil"); - free_stream(g_in_s); - } - make_stream(g_in_s); - init_stream(g_in_s, total_bytes); - g_in_s->end = g_in_s->data + total_bytes; - return audin_data_fragment(chan_id, s); -} - -/*****************************************************************************/ -static int -audin_data(int chan_id, struct stream *s) -{ - LOG_DEVEL_HEXDUMP(LOG_LEVEL_TRACE, "audin_data:", s->p, s_rem(s)); - if (g_in_s == NULL) - { - return audin_process_msg(chan_id, s); - } - return audin_data_fragment(chan_id, s); -} - /*****************************************************************************/ int audin_init(void) @@ -488,10 +432,9 @@ audin_init(void) g_memset(&g_audin_info, 0, sizeof(g_audin_info)); g_audin_info.open_response = audin_open_response; g_audin_info.close_response = audin_close_response; - g_audin_info.data_first = audin_data_first; - g_audin_info.data = audin_data; + g_audin_info.data_first = NULL; + g_audin_info.data = audin_process_msg; g_audin_chanid = 0; - g_in_s = NULL; return 0; } diff --git a/sesman/chansrv/chansrv.c b/sesman/chansrv/chansrv.c index 96ec3e1d..066ddd06 100644 --- a/sesman/chansrv/chansrv.c +++ b/sesman/chansrv/chansrv.c @@ -46,6 +46,7 @@ #include "xrdp_constants.h" #include "audin.h" #include "channel_defs.h" +#include "dechunker.h" #include "scp.h" #include "scp_sync.h" @@ -98,6 +99,7 @@ struct chansrv_drdynvc int status; /* see CHANSRV_DRDYNVC_STATUS_* */ int flags; int pad0; + struct dyn_dechunker *dc; // Use to dechunk fragments int (*open_response)(int chan_id, int creation_status); int (*close_response)(int chan_id); int (*data_first)(int chan_id, struct stream *s, int total_bytes); @@ -105,7 +107,7 @@ struct chansrv_drdynvc struct trans *xrdp_api_trans; }; -static struct chansrv_drdynvc g_drdynvcs[256]; +static struct chansrv_drdynvc g_drdynvcs[DRDYNVC_CHANNEL_COUNT]; /* data in struct trans::callback_data */ struct xrdp_api_data @@ -576,7 +578,7 @@ static int process_message_drdynvc_open_response(struct stream *s) { struct chansrv_drdynvc *drdynvc; - int chan_id; + uint32_t chan_id; int creation_status; LOG_DEVEL(LOG_LEVEL_DEBUG, "process_message_drdynvc_open_response:"); @@ -586,7 +588,7 @@ process_message_drdynvc_open_response(struct stream *s) } in_uint32_le(s, chan_id); in_uint32_le(s, creation_status); - if ((chan_id < 0) || (chan_id > 255)) + if (chan_id >= DRDYNVC_CHANNEL_COUNT) { return 1; } @@ -603,6 +605,8 @@ process_message_drdynvc_open_response(struct stream *s) else { drdynvc->status = CHANSRV_DRDYNVC_STATUS_CLOSED; + dyn_dechunker_free(drdynvc->dc); + drdynvc->dc = NULL; } if (drdynvc->open_response != NULL) { @@ -621,7 +625,7 @@ static int process_message_drdynvc_close_response(struct stream *s) { struct chansrv_drdynvc *drdynvc; - int chan_id; + uint32_t chan_id; LOG_DEVEL(LOG_LEVEL_DEBUG, "process_message_drdynvc_close_response:"); if (!s_check_rem(s, 4)) @@ -629,7 +633,7 @@ process_message_drdynvc_close_response(struct stream *s) return 1; } in_uint32_le(s, chan_id); - if ((chan_id < 0) || (chan_id > 255)) + if (chan_id >= DRDYNVC_CHANNEL_COUNT) { return 1; } @@ -640,6 +644,8 @@ process_message_drdynvc_close_response(struct stream *s) return 0; } drdynvc->status = CHANSRV_DRDYNVC_STATUS_CLOSED; + dyn_dechunker_free(drdynvc->dc); + drdynvc->dc = NULL; if (drdynvc->close_response != NULL) { if (drdynvc->close_response(chan_id) != 0) @@ -659,6 +665,7 @@ process_message_drdynvc_data_first(struct stream *s) struct chansrv_drdynvc *drdynvc; uint32_t chan_id; int total_bytes; + int rv = 0; LOG_DEVEL(LOG_LEVEL_DEBUG, "process_message_drdynvc_data_first:"); if (!s_check_rem(s, 8)) @@ -667,19 +674,47 @@ process_message_drdynvc_data_first(struct stream *s) } in_uint32_le(s, chan_id); in_uint32_le(s, total_bytes); - if (chan_id > 255) + if (chan_id >= DRDYNVC_CHANNEL_COUNT) { return 1; } drdynvc = g_drdynvcs + chan_id; if (drdynvc->data_first != NULL) { - if (drdynvc->data_first(chan_id, s, total_bytes) != 0) + // Caller has requested to defragment PDUs themself + rv = drdynvc->data_first(chan_id, s, total_bytes); + } + else + { + // Get the dechunker working on the stream + enum dyn_dechunker_status status; + status = dyn_dechunker_process_first_chunk(drdynvc->dc, + s, total_bytes); + switch (status) { - return 1; + case E_DYN_INLINE_CHUNK: + // Pass through as data PDU + if (drdynvc->data != NULL) + { + rv = drdynvc->data(chan_id, s); + } + break; + + case E_DYN_IN_PROGRESS: + break; + + case E_DYN_ERROR: + rv = 1; + break; + + default: + LOG(LOG_LEVEL_ERROR, + "Dechunker returned unknown error %d", + (int)status); + rv = 1; } } - return 0; + return rv; } /*****************************************************************************/ @@ -689,23 +724,73 @@ static int process_message_drdynvc_data(struct stream *s) { struct chansrv_drdynvc *drdynvc; - uint32_t chan_id; + int chan_id; + struct stream *ls = NULL; // Set if the application to be called + int free_ls = 0; // Set if we need to clear ls when we're done + int rv = 0; LOG_DEVEL(LOG_LEVEL_DEBUG, "process_message_drdynvc_data:"); - if (!s_check_rem(s, 8)) + if (!s_check_rem(s, 4)) { return 1; } in_uint32_le(s, chan_id); - drdynvc = g_drdynvcs + chan_id; - if (drdynvc->data != NULL) + if (chan_id >= DRDYNVC_CHANNEL_COUNT) { - if (drdynvc->data(chan_id, s) != 0) + return 1; + } + drdynvc = g_drdynvcs + chan_id; + if (drdynvc->data_first != NULL) + { + // Caller is processing all PDUs directly + ls = s; + } + else + { + // Pass the PDU to the dechunker + enum dyn_dechunker_status dechunker_status = + dyn_dechunker_process_data_chunk(drdynvc->dc, s); + switch (dechunker_status) { - return 1; + case E_DYN_INLINE_CHUNK: + ls = s; + break; + + case E_DYN_IN_PROGRESS: + rv = 0; + break; + + case E_DYN_READY: + ls = dyn_dechunker_get_stream(drdynvc->dc); + // We now own the stream, so must delete it + free_ls = 1; + break; + + case E_DYN_ERROR: + // Error has been logged + rv = 1; + break; + + default: + LOG(LOG_LEVEL_ERROR, + "Dechunker returned unknown error %d", + (int)dechunker_status); + rv = 1; } } - return 0; + + if (ls != NULL) + { + if (drdynvc->data != NULL) + { + rv = drdynvc->data(chan_id, ls); + } + if (free_ls) + { + free_stream(ls); + } + } + return rv; } /*****************************************************************************/ @@ -718,12 +803,13 @@ chansrv_drdynvc_open(const char *name, int flags, int name_bytes; int lchan_id; int error; + struct dyn_dechunker *dc = NULL; lchan_id = 1; while (g_drdynvcs[lchan_id].status != CHANSRV_DRDYNVC_STATUS_CLOSED) { lchan_id++; - if (lchan_id > 255) + if (lchan_id >= DRDYNVC_CHANNEL_COUNT) { return 1; } @@ -733,6 +819,18 @@ chansrv_drdynvc_open(const char *name, int flags, { return 1; } + + // If there's no 'data_first' proc, we need a dechunker for the + // channel + if (procs->data_first == NULL) + { + if ((dc = dyn_dechunker_init(name)) == NULL) + { + // Error logged + return 1; + } + } + name_bytes = g_strlen(name); out_uint32_le(s, 0); /* version */ out_uint32_le(s, 8 + 8 + 4 + name_bytes + 4 + 4); @@ -749,13 +847,17 @@ chansrv_drdynvc_open(const char *name, int flags, if (chan_id != NULL) { *chan_id = lchan_id; - g_drdynvcs[lchan_id].open_response = procs->open_response; - g_drdynvcs[lchan_id].close_response = procs->close_response; - g_drdynvcs[lchan_id].data_first = procs->data_first; - g_drdynvcs[lchan_id].data = procs->data; - g_drdynvcs[lchan_id].status = CHANSRV_DRDYNVC_STATUS_OPEN_SENT; - } + g_drdynvcs[lchan_id].open_response = procs->open_response; + g_drdynvcs[lchan_id].close_response = procs->close_response; + g_drdynvcs[lchan_id].data_first = procs->data_first; + g_drdynvcs[lchan_id].data = procs->data; + g_drdynvcs[lchan_id].status = CHANSRV_DRDYNVC_STATUS_OPEN_SENT; + g_drdynvcs[lchan_id].dc = dc; + } + else + { + dyn_dechunker_free(dc); } return error; } @@ -1738,6 +1840,8 @@ x_server_fatal_handler(void) int main_cleanup(void) { + int i; + if (g_term_event != 0) { g_delete_wait_obj(g_term_event); @@ -1756,6 +1860,10 @@ main_cleanup(void) tc_mutex_delete(g_exec_mutex); tc_sem_delete(g_exec_sem); } + for (i = 0 ; i < DRDYNVC_CHANNEL_COUNT; ++i) + { + dyn_dechunker_free(g_drdynvcs[i].dc); + } log_end(); config_free(g_cfg); g_deinit(); /* os_calls */ diff --git a/sesman/chansrv/chansrv.h b/sesman/chansrv/chansrv.h index d085f695..2434cb35 100644 --- a/sesman/chansrv/chansrv.h +++ b/sesman/chansrv/chansrv.h @@ -62,6 +62,8 @@ struct chansrv_drdynvc_procs { int (*open_response)(int chan_id, int creation_status); int (*close_response)(int chan_id); + // Set data_first to NULL to have the dechunker automatically + // handle channel fragments int (*data_first)(int chan_id, struct stream *s, int total_bytes); int (*data)(int chan_id, struct stream *s); };