server.c 43 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105110611071108110911101111111211131114111511161117111811191120112111221123112411251126112711281129113011311132113311341135113611371138113911401141114211431144114511461147114811491150115111521153115411551156115711581159116011611162116311641165116611671168116911701171117211731174117511761177117811791180118111821183118411851186118711881189119011911192119311941195119611971198119912001201120212031204120512061207120812091210121112121213121412151216121712181219122012211222122312241225122612271228122912301231123212331234
  1. /* Gearman server and library
  2. * Copyright (C) 2008 Brian Aker, Eric Day
  3. * All rights reserved.
  4. *
  5. * Use and distribution licensed under the BSD license. See
  6. * the COPYING file in the parent directory for full text.
  7. */
  8. /**
  9. * @file
  10. * @brief Server Definitions
  11. */
  12. #include <libgearman-server/common.h>
  13. #include <assert.h>
  14. #include <errno.h>
  15. #include <string.h>
  16. #include <iso646.h>
  17. /*
  18. * Private declarations
  19. */
  20. #define TEXT_SUCCESS "OK\r\n"
  21. #define TEXT_ERROR_ARGS "ERR INVALID_ARGUMENTS An+incomplete+set+of+arguments+was+sent+to+this+command+%.*s\r\n"
  22. #define TEXT_ERROR_CREATE_FUNCTION "ERR CREATE_FUNCTION %.*s\r\n"
  23. #define TEXT_ERROR_UNKNOWN_COMMAND "ERR UNKNOWN_COMMAND Unknown+server+command%.*s\r\n"
  24. #define TEXT_ERROR_INTERNAL_ERROR "ERR UNKNOWN_ERROR\r\n"
  25. /**
  26. * @addtogroup gearman_server_private Private Server Functions
  27. * @ingroup gearman_server
  28. * @{
  29. */
  30. /**
  31. * Add job to queue wihle replaying queue during startup.
  32. */
  33. static gearmand_error_t _queue_replay_add(gearman_server_st *server, void *context,
  34. const char *unique, size_t unique_size,
  35. const char *function_name, size_t function_name_size,
  36. const void *data, size_t data_size,
  37. gearmand_job_priority_t priority,
  38. int64_t when);
  39. /**
  40. * Queue an error packet.
  41. */
  42. static gearmand_error_t _server_error_packet(gearman_server_con_st *server_con,
  43. const char *error_code,
  44. const char *error_string);
  45. /**
  46. * Process text commands for a connection.
  47. */
  48. static gearmand_error_t _server_run_text(gearman_server_con_st *server_con,
  49. gearmand_packet_st *packet);
  50. /**
  51. * Send work result packets with data back to clients.
  52. */
  53. static gearmand_error_t
  54. _server_queue_work_data(gearman_server_job_st *server_job,
  55. gearmand_packet_st *packet, gearman_command_t command);
  56. /** @} */
  57. /*
  58. * Public definitions
  59. */
  60. gearmand_error_t gearman_server_run_command(gearman_server_con_st *server_con,
  61. gearmand_packet_st *packet)
  62. {
  63. gearmand_error_t ret;
  64. gearman_server_job_st *server_job;
  65. char job_handle[GEARMAND_JOB_HANDLE_SIZE];
  66. char option[GEARMAN_OPTION_SIZE];
  67. gearman_server_client_st *server_client= NULL;
  68. char numerator_buffer[11]; /* Max string size to hold a uint32_t. */
  69. char denominator_buffer[11]; /* Max string size to hold a uint32_t. */
  70. gearmand_job_priority_t priority;
  71. int checked_length;
  72. if (packet->magic == GEARMAN_MAGIC_RESPONSE)
  73. {
  74. return _server_error_packet(server_con, "bad_magic", "Request magic expected");
  75. }
  76. gearmand_log_debug(GEARMAN_DEFAULT_LOG_PARAM,
  77. "%15s:%5s packet command %s",
  78. server_con->con.context == NULL ? "-" : server_con->con.context->host,
  79. server_con->con.context == NULL ? "-" : server_con->con.context->port,
  80. gearmand_strcommand(packet));
  81. switch (packet->command)
  82. {
  83. /* Client/worker requests. */
  84. case GEARMAN_COMMAND_ECHO_REQ:
  85. /* Reuse the data buffer and just shove the data back. */
  86. ret= gearman_server_io_packet_add(server_con, true, GEARMAN_MAGIC_RESPONSE,
  87. GEARMAN_COMMAND_ECHO_RES, packet->data,
  88. packet->data_size, NULL);
  89. if (gearmand_failed(ret))
  90. {
  91. return gearmand_gerror("gearman_server_io_packet_add", ret);
  92. }
  93. packet->options.free_data= false;
  94. break;
  95. case GEARMAN_COMMAND_SUBMIT_REDUCE_JOB: // Reduce request
  96. server_client= gearman_server_client_add(server_con);
  97. if (server_client == NULL)
  98. {
  99. return GEARMAN_MEMORY_ALLOCATION_FAILURE;
  100. }
  101. case GEARMAN_COMMAND_SUBMIT_REDUCE_JOB_BACKGROUND:
  102. {
  103. gearmand_log_debug(GEARMAN_DEFAULT_LOG_PARAM,
  104. "Received reduce submission, Partitioner: %.*s(%lu) Reducer: %.*s(%lu) Unique: %.*s(%lu) with %d arguments",
  105. packet->arg_size[0] -1, packet->arg[0], packet->arg_size[0] -1,
  106. packet->arg_size[2] -1, packet->arg[2], packet->arg_size[2] -1, // reducer
  107. packet->arg_size[1] -1, packet->arg[1], packet->arg_size[1] -1,
  108. (int)packet->argc);
  109. if (packet->arg_size[2] -1 > GEARMAN_UNIQUE_SIZE)
  110. {
  111. gearman_server_client_free(server_client);
  112. gearmand_gerror("unique value too large", GEARMAN_ARGUMENT_TOO_LARGE);
  113. return _server_error_packet(server_con, "job failure", "Unique value too large");
  114. }
  115. gearmand_job_priority_t map_priority= GEARMAND_JOB_PRIORITY_NORMAL;
  116. /* Schedule job. */
  117. server_job= gearman_server_job_add_reducer(Server,
  118. (char *)(packet->arg[0]), packet->arg_size[0] -1, // Function
  119. (char *)(packet->arg[1]), packet->arg_size[1] -1, // unique
  120. (char *)(packet->arg[2]), packet->arg_size[2] -1, // reducer
  121. packet->data, packet->data_size, map_priority,
  122. server_client, &ret, 0);
  123. if (gearmand_success(ret))
  124. {
  125. packet->options.free_data= false;
  126. }
  127. else if (ret == GEARMAN_JOB_QUEUE_FULL)
  128. {
  129. gearman_server_client_free(server_client);
  130. return _server_error_packet(server_con, "queue_full", "Job queue is full");
  131. }
  132. else if (ret != GEARMAN_JOB_EXISTS)
  133. {
  134. gearman_server_client_free(server_client);
  135. gearmand_gerror("gearman_server_job_add", ret);
  136. return ret;
  137. }
  138. /* Queue the job created packet. */
  139. ret= gearman_server_io_packet_add(server_con, false, GEARMAN_MAGIC_RESPONSE,
  140. GEARMAN_COMMAND_JOB_CREATED,
  141. server_job->job_handle,
  142. (size_t)strlen(server_job->job_handle),
  143. NULL);
  144. if (gearmand_failed(ret))
  145. {
  146. gearman_server_client_free(server_client);
  147. gearmand_gerror("gearman_server_io_packet_add", ret);
  148. return ret;
  149. }
  150. gearmand_log_notice(GEARMAN_DEFAULT_LOG_PARAM,"accepted,%.*s,%.*s,%.*s",
  151. packet->arg_size[0] -1, packet->arg[0], // Function
  152. packet->arg_size[1] -1, packet->arg[1], // unique
  153. packet->arg_size[2] -1, packet->arg[2]); // reducer
  154. }
  155. break;
  156. /* Client requests. */
  157. case GEARMAN_COMMAND_SUBMIT_JOB:
  158. case GEARMAN_COMMAND_SUBMIT_JOB_BG:
  159. case GEARMAN_COMMAND_SUBMIT_JOB_HIGH:
  160. case GEARMAN_COMMAND_SUBMIT_JOB_HIGH_BG:
  161. case GEARMAN_COMMAND_SUBMIT_JOB_LOW:
  162. case GEARMAN_COMMAND_SUBMIT_JOB_LOW_BG:
  163. case GEARMAN_COMMAND_SUBMIT_JOB_EPOCH:
  164. if (packet->command == GEARMAN_COMMAND_SUBMIT_JOB ||
  165. packet->command == GEARMAN_COMMAND_SUBMIT_JOB_BG ||
  166. packet->command == GEARMAN_COMMAND_SUBMIT_JOB_EPOCH)
  167. {
  168. priority= GEARMAND_JOB_PRIORITY_NORMAL;
  169. }
  170. else if (packet->command == GEARMAN_COMMAND_SUBMIT_JOB_HIGH ||
  171. packet->command == GEARMAN_COMMAND_SUBMIT_JOB_HIGH_BG)
  172. {
  173. priority= GEARMAND_JOB_PRIORITY_HIGH;
  174. }
  175. else
  176. {
  177. priority= GEARMAND_JOB_PRIORITY_LOW;
  178. }
  179. if (packet->command == GEARMAN_COMMAND_SUBMIT_JOB_BG ||
  180. packet->command == GEARMAN_COMMAND_SUBMIT_JOB_HIGH_BG ||
  181. packet->command == GEARMAN_COMMAND_SUBMIT_JOB_LOW_BG ||
  182. packet->command == GEARMAN_COMMAND_SUBMIT_JOB_EPOCH)
  183. {
  184. server_client= NULL;
  185. }
  186. else
  187. {
  188. server_client= gearman_server_client_add(server_con);
  189. if (server_client == NULL)
  190. {
  191. return GEARMAN_MEMORY_ALLOCATION_FAILURE;
  192. }
  193. }
  194. gearmand_log_debug(GEARMAN_DEFAULT_LOG_PARAM,
  195. "Received submission, %.*s/%.*s with %d arguments",
  196. packet->arg_size[0], packet->arg[0],
  197. packet->arg_size[1], packet->arg[1],
  198. (int)packet->argc);
  199. int64_t when= 0;
  200. if (packet->command == GEARMAN_COMMAND_SUBMIT_JOB_EPOCH)
  201. {
  202. sscanf((char *)packet->arg[2], "%lld", (long long *)&when);
  203. gearmand_log_debug(GEARMAN_DEFAULT_LOG_PARAM,
  204. "Received EPOCH job submission, %.*s/%.*s, with data for %jd at %jd, args %d",
  205. packet->arg_size[0], packet->arg[0],
  206. packet->arg_size[1], packet->arg[1],
  207. when, time(NULL),
  208. (int)packet->argc);
  209. }
  210. if (packet->arg_size[1] -1 > GEARMAN_UNIQUE_SIZE)
  211. {
  212. gearmand_gerror("unique value too large", GEARMAN_ARGUMENT_TOO_LARGE);
  213. gearman_server_client_free(server_client);
  214. return _server_error_packet(server_con, "job failure", "Unique value too large");
  215. }
  216. /* Schedule job. */
  217. server_job= gearman_server_job_add(Server,
  218. (char *)(packet->arg[0]), packet->arg_size[0] -1, // Function
  219. (char *)(packet->arg[1]), packet->arg_size[1] -1, // unique
  220. packet->data, packet->data_size, priority,
  221. server_client, &ret,
  222. when);
  223. if (gearmand_success(ret))
  224. {
  225. packet->options.free_data= false;
  226. }
  227. else if (ret == GEARMAN_JOB_QUEUE_FULL)
  228. {
  229. gearman_server_client_free(server_client);
  230. return _server_error_packet(server_con, "queue_full", "Job queue is full");
  231. }
  232. else if (ret != GEARMAN_JOB_EXISTS)
  233. {
  234. gearman_server_client_free(server_client);
  235. gearmand_gerror("gearman_server_job_add", ret);
  236. return ret;
  237. }
  238. /* Queue the job created packet. */
  239. ret= gearman_server_io_packet_add(server_con, false, GEARMAN_MAGIC_RESPONSE,
  240. GEARMAN_COMMAND_JOB_CREATED,
  241. server_job->job_handle,
  242. (size_t)strlen(server_job->job_handle),
  243. NULL);
  244. if (gearmand_failed(ret))
  245. {
  246. gearman_server_client_free(server_client);
  247. gearmand_gerror("gearman_server_io_packet_add", ret);
  248. return ret;
  249. }
  250. gearmand_log_notice(GEARMAN_DEFAULT_LOG_PARAM,"accepted,%.*s,%.*s,%jd",
  251. packet->arg_size[0], packet->arg[0], // Function
  252. packet->arg_size[1], packet->arg[1], // Unique
  253. when);
  254. break;
  255. case GEARMAN_COMMAND_GET_STATUS:
  256. {
  257. /* This may not be NULL terminated, so copy to make sure it is. */
  258. checked_length= snprintf(job_handle, GEARMAND_JOB_HANDLE_SIZE, "%.*s",
  259. (int)(packet->arg_size[0]), (char *)(packet->arg[0]));
  260. if (checked_length >= GEARMAND_JOB_HANDLE_SIZE || checked_length < 0)
  261. {
  262. gearmand_log_error(GEARMAN_DEFAULT_LOG_PARAM, "snprintf(%d)", checked_length);
  263. return GEARMAN_MEMORY_ALLOCATION_FAILURE;
  264. }
  265. server_job= gearman_server_job_get(Server, job_handle, NULL);
  266. /* Queue status result packet. */
  267. if (server_job == NULL)
  268. {
  269. ret= gearman_server_io_packet_add(server_con, false,
  270. GEARMAN_MAGIC_RESPONSE,
  271. GEARMAN_COMMAND_STATUS_RES, job_handle,
  272. (size_t)(strlen(job_handle) + 1),
  273. "0", (size_t)2, "0", (size_t)2, "0",
  274. (size_t)2, "0", (size_t)1, NULL);
  275. }
  276. else
  277. {
  278. checked_length= snprintf(numerator_buffer, sizeof(numerator_buffer), "%u", server_job->numerator);
  279. if ((size_t)checked_length >= sizeof(numerator_buffer) || checked_length < 0)
  280. {
  281. gearmand_log_error(GEARMAN_DEFAULT_LOG_PARAM, "snprintf(%d)", checked_length);
  282. return GEARMAN_MEMORY_ALLOCATION_FAILURE;
  283. }
  284. checked_length= snprintf(denominator_buffer, sizeof(denominator_buffer), "%u", server_job->denominator);
  285. if ((size_t)checked_length >= sizeof(denominator_buffer) || checked_length < 0)
  286. {
  287. gearmand_log_error(GEARMAN_DEFAULT_LOG_PARAM, "snprintf(%d)", checked_length);
  288. return GEARMAN_MEMORY_ALLOCATION_FAILURE;
  289. }
  290. ret= gearman_server_io_packet_add(server_con, false,
  291. GEARMAN_MAGIC_RESPONSE,
  292. GEARMAN_COMMAND_STATUS_RES, job_handle,
  293. (size_t)(strlen(job_handle) + 1),
  294. "1", (size_t)2,
  295. server_job->worker == NULL ? "0" : "1",
  296. (size_t)2, numerator_buffer,
  297. (size_t)(strlen(numerator_buffer) + 1),
  298. denominator_buffer,
  299. (size_t)strlen(denominator_buffer),
  300. NULL);
  301. }
  302. if (ret != GEARMAN_SUCCESS)
  303. {
  304. gearmand_gerror("gearman_server_io_packet_add", ret);
  305. return ret;
  306. }
  307. }
  308. break;
  309. case GEARMAN_COMMAND_OPTION_REQ:
  310. /* This may not be NULL terminated, so copy to make sure it is. */
  311. checked_length= snprintf(option, GEARMAN_OPTION_SIZE, "%.*s",
  312. (int)(packet->arg_size[0]), (char *)(packet->arg[0]));
  313. if (checked_length >= GEARMAN_OPTION_SIZE || checked_length < 0)
  314. {
  315. gearmand_log_error(GEARMAN_DEFAULT_LOG_PARAM, "snprintf(%d)", checked_length);
  316. return _server_error_packet(server_con, "unknown_option",
  317. "Server does not recognize given option");
  318. }
  319. if (strcasecmp(option, "exceptions") == 0)
  320. {
  321. gearmand_log_debug(GEARMAN_DEFAULT_LOG_PARAM, "'exceptions'");
  322. server_con->is_exceptions= true;
  323. }
  324. else
  325. {
  326. return _server_error_packet(server_con, "unknown_option",
  327. "Server does not recognize given option");
  328. }
  329. ret= gearman_server_io_packet_add(server_con, false, GEARMAN_MAGIC_RESPONSE,
  330. GEARMAN_COMMAND_OPTION_RES,
  331. packet->arg[0], packet->arg_size[0],
  332. NULL);
  333. if (ret != GEARMAN_SUCCESS)
  334. {
  335. gearmand_gerror("gearman_server_io_packet_add", ret);
  336. return ret;
  337. }
  338. break;
  339. /* Worker requests. */
  340. case GEARMAN_COMMAND_CAN_DO:
  341. if (gearman_server_worker_add(server_con, (char *)(packet->arg[0]),
  342. packet->arg_size[0], 0) == NULL)
  343. {
  344. return GEARMAN_MEMORY_ALLOCATION_FAILURE;
  345. }
  346. break;
  347. case GEARMAN_COMMAND_CAN_DO_TIMEOUT:
  348. if (gearman_server_worker_add(server_con, (char *)(packet->arg[0]),
  349. packet->arg_size[0] - 1,
  350. (in_port_t)atoi((char *)(packet->arg[1])))
  351. == NULL)
  352. {
  353. return GEARMAN_MEMORY_ALLOCATION_FAILURE;
  354. }
  355. break;
  356. case GEARMAN_COMMAND_CANT_DO:
  357. gearman_server_con_free_worker(server_con, (char *)(packet->arg[0]),
  358. packet->arg_size[0]);
  359. break;
  360. case GEARMAN_COMMAND_RESET_ABILITIES:
  361. gearman_server_con_free_workers(server_con);
  362. break;
  363. case GEARMAN_COMMAND_PRE_SLEEP:
  364. server_job= gearman_server_job_peek(server_con);
  365. if (server_job == NULL)
  366. {
  367. server_con->is_sleeping= true;
  368. }
  369. else
  370. {
  371. /* If there are jobs that could be run, queue a NOOP packet to wake the
  372. worker up. This could be the result of a race codition. */
  373. ret= gearman_server_io_packet_add(server_con, false,
  374. GEARMAN_MAGIC_RESPONSE,
  375. GEARMAN_COMMAND_NOOP, NULL);
  376. if (ret != GEARMAN_SUCCESS)
  377. {
  378. gearmand_gerror("gearman_server_io_packet_add", ret);
  379. return ret;
  380. }
  381. }
  382. break;
  383. case GEARMAN_COMMAND_GRAB_JOB:
  384. case GEARMAN_COMMAND_GRAB_JOB_UNIQ:
  385. case GEARMAN_COMMAND_GRAB_JOB_ALL:
  386. server_con->is_sleeping= false;
  387. server_con->is_noop_sent= false;
  388. server_job= gearman_server_job_take(server_con);
  389. if (server_job == NULL)
  390. {
  391. /* No jobs found, queue no job packet. */
  392. ret= gearman_server_io_packet_add(server_con, false,
  393. GEARMAN_MAGIC_RESPONSE,
  394. GEARMAN_COMMAND_NO_JOB, NULL);
  395. }
  396. else if (packet->command == GEARMAN_COMMAND_GRAB_JOB_UNIQ)
  397. {
  398. /*
  399. We found a runnable job, queue job assigned packet and take the job off the queue.
  400. */
  401. ret= gearman_server_io_packet_add(server_con, false,
  402. GEARMAN_MAGIC_RESPONSE,
  403. GEARMAN_COMMAND_JOB_ASSIGN_UNIQ,
  404. server_job->job_handle, (size_t)(strlen(server_job->job_handle) + 1),
  405. server_job->function->function_name, server_job->function->function_name_size + 1,
  406. server_job->unique, (size_t)(strlen(server_job->unique) + 1),
  407. server_job->data, server_job->data_size,
  408. NULL);
  409. }
  410. else if (packet->command == GEARMAN_COMMAND_GRAB_JOB_ALL and server_job->reducer)
  411. {
  412. gearmand_log_debug(GEARMAN_DEFAULT_LOG_PARAM,
  413. "Sending reduce submission, Partitioner: %.*s(%lu) Reducer: %.*s(%lu) Unique: %.*s(%lu) with data sized (%lu)" ,
  414. server_job->function->function_name_size, server_job->function->function_name, server_job->function->function_name_size,
  415. strlen(server_job->reducer), server_job->reducer, strlen(server_job->reducer),
  416. strlen(server_job->unique), server_job->unique, strlen(server_job->unique),
  417. (unsigned long)server_job->data_size);
  418. /*
  419. We found a runnable job, queue job assigned packet and take the job off the queue.
  420. */
  421. ret= gearman_server_io_packet_add(server_con, false,
  422. GEARMAN_MAGIC_RESPONSE,
  423. GEARMAN_COMMAND_JOB_ASSIGN_ALL,
  424. server_job->job_handle, (size_t)(strlen(server_job->job_handle) + 1),
  425. server_job->function->function_name, server_job->function->function_name_size + 1,
  426. server_job->unique, (size_t)(strlen(server_job->unique) + 1),
  427. server_job->reducer, (size_t)(strlen(server_job->reducer) +1),
  428. server_job->data, server_job->data_size,
  429. NULL);
  430. }
  431. else if (packet->command == GEARMAN_COMMAND_GRAB_JOB_ALL)
  432. {
  433. /*
  434. We found a runnable job, queue job assigned packet and take the job off the queue.
  435. */
  436. ret= gearman_server_io_packet_add(server_con, false,
  437. GEARMAN_MAGIC_RESPONSE,
  438. GEARMAN_COMMAND_JOB_ASSIGN_UNIQ,
  439. server_job->job_handle, (size_t)(strlen(server_job->job_handle) +1),
  440. server_job->function->function_name, server_job->function->function_name_size +1,
  441. server_job->unique, (size_t)(strlen(server_job->unique) +1),
  442. server_job->data, server_job->data_size,
  443. NULL);
  444. }
  445. else
  446. {
  447. /* Same, but without unique ID. */
  448. ret= gearman_server_io_packet_add(server_con, false,
  449. GEARMAN_MAGIC_RESPONSE,
  450. GEARMAN_COMMAND_JOB_ASSIGN,
  451. server_job->job_handle, (size_t)(strlen(server_job->job_handle) + 1),
  452. server_job->function->function_name, server_job->function->function_name_size + 1,
  453. server_job->data, server_job->data_size,
  454. NULL);
  455. }
  456. if (gearmand_failed(ret))
  457. {
  458. gearmand_gerror("gearman_server_io_packet_add", ret);
  459. if (server_job)
  460. return gearman_server_job_queue(server_job);
  461. return ret;
  462. }
  463. break;
  464. case GEARMAN_COMMAND_WORK_DATA:
  465. case GEARMAN_COMMAND_WORK_WARNING:
  466. server_job= gearman_server_job_get(Server,
  467. (char *)(packet->arg[0]),
  468. server_con);
  469. if (server_job == NULL)
  470. {
  471. return _server_error_packet(server_con, "job_not_found",
  472. "Job given in work result not found");
  473. }
  474. /* Queue the data/warning packet for all clients. */
  475. ret= _server_queue_work_data(server_job, packet, packet->command);
  476. if (gearmand_failed(ret))
  477. {
  478. gearmand_gerror("_server_queue_work_data", ret);
  479. return ret;
  480. }
  481. break;
  482. case GEARMAN_COMMAND_WORK_STATUS:
  483. server_job= gearman_server_job_get(Server,
  484. (char *)(packet->arg[0]),
  485. server_con);
  486. if (server_job == NULL)
  487. {
  488. return _server_error_packet(server_con, "job_not_found",
  489. "Job given in work result not found");
  490. }
  491. /* Update job status. */
  492. server_job->numerator= (uint32_t)atoi((char *)(packet->arg[1]));
  493. /* This may not be NULL terminated, so copy to make sure it is. */
  494. checked_length= snprintf(denominator_buffer, sizeof(denominator_buffer), "%.*s", (int)(packet->arg_size[2]),
  495. (char *)(packet->arg[2]));
  496. if ((size_t)checked_length > sizeof(denominator_buffer) || checked_length < 0)
  497. {
  498. gearmand_error("snprintf");
  499. return GEARMAN_MEMORY_ALLOCATION_FAILURE;
  500. }
  501. server_job->denominator= (uint32_t)atoi(denominator_buffer);
  502. /* Queue the status packet for all clients. */
  503. for (server_client= server_job->client_list; server_client;
  504. server_client= server_client->job_next)
  505. {
  506. ret= gearman_server_io_packet_add(server_client->con, false,
  507. GEARMAN_MAGIC_RESPONSE,
  508. GEARMAN_COMMAND_WORK_STATUS,
  509. packet->arg[0], packet->arg_size[0],
  510. packet->arg[1], packet->arg_size[1],
  511. packet->arg[2], packet->arg_size[2],
  512. NULL);
  513. if (gearmand_failed(ret))
  514. {
  515. gearmand_gerror("gearman_server_io_packet_add", ret);
  516. return ret;
  517. }
  518. }
  519. break;
  520. case GEARMAN_COMMAND_WORK_COMPLETE:
  521. server_job= gearman_server_job_get(Server,
  522. (char *)(packet->arg[0]),
  523. server_con);
  524. if (server_job == NULL)
  525. {
  526. return _server_error_packet(server_con, "job_not_found", "Job given in work result not found");
  527. }
  528. /* Queue the complete packet for all clients. */
  529. ret= _server_queue_work_data(server_job, packet,
  530. GEARMAN_COMMAND_WORK_COMPLETE);
  531. if (gearmand_failed(ret))
  532. {
  533. gearmand_gerror("_server_queue_work_data", ret);
  534. return ret;
  535. }
  536. /* Remove from persistent queue if one exists. */
  537. if (server_job->job_queued && Server->queue._done_fn != NULL)
  538. {
  539. ret= (*(Server->queue._done_fn))(Server, (void *)Server->queue._context,
  540. server_job->unique,
  541. (size_t)strlen(server_job->unique),
  542. server_job->function->function_name,
  543. server_job->function->function_name_size);
  544. if (gearmand_failed(ret))
  545. {
  546. gearmand_gerror("Remove from persistent queue", ret);
  547. return ret;
  548. }
  549. }
  550. /* Job is done, remove it. */
  551. gearman_server_job_free(server_job);
  552. break;
  553. case GEARMAN_COMMAND_WORK_EXCEPTION:
  554. server_job= gearman_server_job_get(Server,
  555. (char *)(packet->arg[0]),
  556. server_con);
  557. if (server_job == NULL)
  558. {
  559. return _server_error_packet(server_con, "job_not_found", "Job given in work result not found");
  560. }
  561. /* Queue the exception packet for all clients. */
  562. ret= _server_queue_work_data(server_job, packet,
  563. GEARMAN_COMMAND_WORK_EXCEPTION);
  564. if (gearmand_failed(ret))
  565. {
  566. gearmand_gerror("_server_queue_work_data", ret);
  567. return ret;
  568. }
  569. break;
  570. case GEARMAN_COMMAND_WORK_FAIL:
  571. /* This may not be NULL terminated, so copy to make sure it is. */
  572. checked_length= snprintf(job_handle, GEARMAND_JOB_HANDLE_SIZE, "%.*s",
  573. (int)(packet->arg_size[0]), (char *)(packet->arg[0]));
  574. if (checked_length >= GEARMAND_JOB_HANDLE_SIZE || checked_length < 0)
  575. {
  576. return _server_error_packet(server_con, "job_name_too_large",
  577. "Error occured due to GEARMAND_JOB_HANDLE_SIZE being too small from snprintf");
  578. }
  579. server_job= gearman_server_job_get(Server, job_handle,
  580. server_con);
  581. if (server_job == NULL)
  582. {
  583. return _server_error_packet(server_con, "job_not_found",
  584. "Job given in work result not found");
  585. }
  586. /* Queue the fail packet for all clients. */
  587. for (server_client= server_job->client_list; server_client;
  588. server_client= server_client->job_next)
  589. {
  590. ret= gearman_server_io_packet_add(server_client->con, false,
  591. GEARMAN_MAGIC_RESPONSE,
  592. GEARMAN_COMMAND_WORK_FAIL,
  593. packet->arg[0], packet->arg_size[0],
  594. NULL);
  595. if (gearmand_failed(ret))
  596. {
  597. gearmand_gerror("gearman_server_io_packet_add", ret);
  598. return ret;
  599. }
  600. }
  601. /* Remove from persistent queue if one exists. */
  602. if (server_job->job_queued && Server->queue._done_fn != NULL)
  603. {
  604. ret= (*(Server->queue._done_fn))(Server, (void *)Server->queue._context,
  605. server_job->unique,
  606. (size_t)strlen(server_job->unique),
  607. server_job->function->function_name,
  608. server_job->function->function_name_size);
  609. if (gearmand_failed(ret))
  610. {
  611. gearmand_gerror("Remove from persistent queue", ret);
  612. return ret;
  613. }
  614. }
  615. /* Job is done, remove it. */
  616. gearman_server_job_free(server_job);
  617. break;
  618. case GEARMAN_COMMAND_SET_CLIENT_ID:
  619. gearman_server_con_set_id(server_con, (char *)(packet->arg[0]),
  620. packet->arg_size[0]);
  621. break;
  622. case GEARMAN_COMMAND_TEXT:
  623. return _server_run_text(server_con, packet);
  624. case GEARMAN_COMMAND_UNUSED:
  625. case GEARMAN_COMMAND_NOOP:
  626. case GEARMAN_COMMAND_JOB_CREATED:
  627. case GEARMAN_COMMAND_NO_JOB:
  628. case GEARMAN_COMMAND_JOB_ASSIGN:
  629. case GEARMAN_COMMAND_ECHO_RES:
  630. case GEARMAN_COMMAND_ERROR:
  631. case GEARMAN_COMMAND_STATUS_RES:
  632. case GEARMAN_COMMAND_ALL_YOURS:
  633. case GEARMAN_COMMAND_OPTION_RES:
  634. case GEARMAN_COMMAND_SUBMIT_JOB_SCHED:
  635. case GEARMAN_COMMAND_JOB_ASSIGN_UNIQ:
  636. case GEARMAN_COMMAND_JOB_ASSIGN_ALL:
  637. case GEARMAN_COMMAND_MAX:
  638. default:
  639. return _server_error_packet(server_con, "bad_command", "Command not expected");
  640. }
  641. return GEARMAN_SUCCESS;
  642. }
  643. gearmand_error_t gearman_server_shutdown_graceful(gearman_server_st *server)
  644. {
  645. server->shutdown_graceful= true;
  646. if (server->job_count == 0)
  647. return GEARMAN_SHUTDOWN;
  648. return GEARMAN_SHUTDOWN_GRACEFUL;
  649. }
  650. gearmand_error_t gearman_server_queue_replay(gearman_server_st *server)
  651. {
  652. gearmand_error_t ret;
  653. if (server->queue._replay_fn == NULL)
  654. return GEARMAN_SUCCESS;
  655. server->state.queue_startup= true;
  656. ret= (*(server->queue._replay_fn))(server, (void *)server->queue._context,
  657. _queue_replay_add, server);
  658. server->state.queue_startup= false;
  659. return ret;
  660. }
  661. void *gearman_server_queue_context(const gearman_server_st *server)
  662. {
  663. return (void *)server->queue._context;
  664. }
  665. void gearman_server_set_queue(gearman_server_st *server,
  666. void *context,
  667. gearman_queue_add_fn *add,
  668. gearman_queue_flush_fn *flush,
  669. gearman_queue_done_fn *done,
  670. gearman_queue_replay_fn *replay)
  671. {
  672. server->queue._context= context;
  673. server->queue._add_fn= add;
  674. server->queue._flush_fn= flush;
  675. server->queue._done_fn= done;
  676. server->queue._replay_fn= replay;
  677. }
  678. /*
  679. * Private definitions
  680. */
  681. gearmand_error_t _queue_replay_add(gearman_server_st *server,
  682. void *context __attribute__ ((unused)),
  683. const char *unique, size_t unique_size,
  684. const char *function_name, size_t function_name_size,
  685. const void *data, size_t data_size,
  686. gearmand_job_priority_t priority,
  687. int64_t when)
  688. {
  689. gearmand_error_t ret= GEARMAN_SUCCESS;
  690. (void)gearman_server_job_add(server,
  691. function_name, function_name_size,
  692. unique, unique_size,
  693. data, data_size, priority, NULL, &ret, when);
  694. if (ret != GEARMAN_SUCCESS)
  695. gearmand_gerror("gearman_server_job_add", ret);
  696. return ret;
  697. }
  698. static gearmand_error_t _server_error_packet(gearman_server_con_st *server_con,
  699. const char *error_code,
  700. const char *error_string)
  701. {
  702. return gearman_server_io_packet_add(server_con, false, GEARMAN_MAGIC_RESPONSE,
  703. GEARMAN_COMMAND_ERROR, error_code,
  704. (size_t)(strlen(error_code) + 1),
  705. error_string,
  706. (size_t)strlen(error_string), NULL);
  707. }
  708. static gearmand_error_t _server_run_text(gearman_server_con_st *server_con,
  709. gearmand_packet_st *packet)
  710. {
  711. char *data;
  712. char *new_data;
  713. size_t size;
  714. size_t total;
  715. int max_queue_size;
  716. gearman_server_thread_st *thread;
  717. gearman_server_con_st *con;
  718. gearman_server_worker_st *worker;
  719. gearman_server_packet_st *server_packet;
  720. int checked_length;
  721. data= (char *)(char *)malloc(GEARMAN_TEXT_RESPONSE_SIZE);
  722. if (data == NULL)
  723. {
  724. gearmand_perror("malloc");
  725. return GEARMAN_MEMORY_ALLOCATION_FAILURE;
  726. }
  727. total= GEARMAN_TEXT_RESPONSE_SIZE;
  728. data[0]= 0;
  729. if (packet->argc)
  730. {
  731. gearmand_log_debug(GEARMAN_DEFAULT_LOG_PARAM, "text command %.*s", packet->arg_size[0], packet->arg[0]);
  732. }
  733. if (packet->argc == 0)
  734. {
  735. snprintf(data, GEARMAN_TEXT_RESPONSE_SIZE, TEXT_ERROR_UNKNOWN_COMMAND, 4, "NULL");
  736. }
  737. else if (!strcasecmp("workers", (char *)(packet->arg[0])))
  738. {
  739. size= 0;
  740. for (thread= Server->thread_list; thread != NULL;
  741. thread= thread->next)
  742. {
  743. int error;
  744. if (! (error= pthread_mutex_lock(&thread->lock)))
  745. {
  746. for (con= thread->con_list; con != NULL; con= con->next)
  747. {
  748. if (con->_host == NULL)
  749. continue;
  750. if (size > total)
  751. size= total;
  752. /* Make sure we have at least GEARMAN_TEXT_RESPONSE_SIZE bytes. */
  753. if (size + GEARMAN_TEXT_RESPONSE_SIZE > total)
  754. {
  755. new_data= (char *)realloc(data, total + GEARMAN_TEXT_RESPONSE_SIZE);
  756. if (new_data == NULL)
  757. {
  758. (void) pthread_mutex_unlock(&thread->lock);
  759. gearmand_perror("realloc");
  760. gearmand_debug("free");
  761. free(data);
  762. return GEARMAN_MEMORY_ALLOCATION_FAILURE;
  763. }
  764. data= new_data;
  765. total+= GEARMAN_TEXT_RESPONSE_SIZE;
  766. }
  767. checked_length= snprintf(data + size, total - size, "%d %s %s :",
  768. con->con.fd, con->_host, con->id);
  769. if ((size_t)checked_length > total - size || checked_length < 0)
  770. {
  771. (void) pthread_mutex_unlock(&thread->lock);
  772. gearmand_debug("free");
  773. free(data);
  774. gearmand_perror("snprintf");
  775. return GEARMAN_MEMORY_ALLOCATION_FAILURE;
  776. }
  777. size+= (size_t)checked_length;
  778. if (size > total)
  779. continue;
  780. for (worker= con->worker_list; worker != NULL; worker= worker->con_next)
  781. {
  782. checked_length= snprintf(data + size, total - size, " %.*s",
  783. (int)(worker->function->function_name_size),
  784. worker->function->function_name);
  785. if ((size_t)checked_length > total - size || checked_length < 0)
  786. {
  787. (void) pthread_mutex_unlock(&thread->lock);
  788. gearmand_debug("free");
  789. free(data);
  790. gearmand_perror("snprintf");
  791. return GEARMAN_MEMORY_ALLOCATION_FAILURE;
  792. }
  793. size+= (size_t)checked_length;
  794. if (size > total)
  795. break;
  796. }
  797. if (size > total)
  798. continue;
  799. checked_length= snprintf(data + size, total - size, "\n");
  800. if ((size_t)checked_length > total - size || checked_length < 0)
  801. {
  802. (void) pthread_mutex_unlock(&thread->lock);
  803. gearmand_debug("free");
  804. free(data);
  805. gearmand_perror("snprintf");
  806. return GEARMAN_MEMORY_ALLOCATION_FAILURE;
  807. }
  808. size+= (size_t)checked_length;
  809. }
  810. (void) pthread_mutex_unlock(&thread->lock);
  811. }
  812. else
  813. {
  814. errno= error;
  815. gearmand_error("pthread_mutex_lock");
  816. assert(! "pthread_mutex_lock");
  817. }
  818. }
  819. if (size < total)
  820. {
  821. checked_length= snprintf(data + size, total - size, ".\n");
  822. if ((size_t)checked_length > total - size || checked_length < 0)
  823. {
  824. gearmand_perror("snprintf");
  825. return GEARMAN_MEMORY_ALLOCATION_FAILURE;
  826. }
  827. }
  828. }
  829. else if (!strcasecmp("status", (char *)(packet->arg[0])))
  830. {
  831. size= 0;
  832. gearman_server_function_st *function;
  833. for (function= Server->function_list; function != NULL;
  834. function= function->next)
  835. {
  836. if (size + GEARMAN_TEXT_RESPONSE_SIZE > total)
  837. {
  838. new_data= (char *)realloc(data, total + GEARMAN_TEXT_RESPONSE_SIZE);
  839. if (new_data == NULL)
  840. {
  841. gearmand_perror("realloc");
  842. gearmand_debug("free");
  843. free(data);
  844. return GEARMAN_MEMORY_ALLOCATION_FAILURE;
  845. }
  846. data= new_data;
  847. total+= GEARMAN_TEXT_RESPONSE_SIZE;
  848. }
  849. checked_length= snprintf(data + size, total - size, "%.*s\t%u\t%u\t%u\n",
  850. (int)(function->function_name_size),
  851. function->function_name, function->job_total,
  852. function->job_running, function->worker_count);
  853. if ((size_t)checked_length > total - size || checked_length < 0)
  854. {
  855. gearmand_perror("snprintf");
  856. gearmand_debug("free");
  857. free(data);
  858. return GEARMAN_MEMORY_ALLOCATION_FAILURE;
  859. }
  860. size+= (size_t)checked_length;
  861. if (size > total)
  862. size= total;
  863. }
  864. if (size < total)
  865. {
  866. checked_length= snprintf(data + size, total - size, ".\n");
  867. if ((size_t)checked_length > total - size || checked_length < 0)
  868. {
  869. gearmand_perror("snprintf");
  870. gearmand_debug("free");
  871. free(data);
  872. return GEARMAN_MEMORY_ALLOCATION_FAILURE;
  873. }
  874. }
  875. }
  876. else if (!strcasecmp("create", (char *)(packet->arg[0])))
  877. {
  878. if (packet->argc == 3 && !strcasecmp("function", (char *)(packet->arg[1])))
  879. {
  880. gearman_server_function_st *function;
  881. function= gearman_server_function_get(Server, (char *)(packet->arg[2]), packet->arg_size[2] -2);
  882. if (function)
  883. {
  884. snprintf(data, GEARMAN_TEXT_RESPONSE_SIZE, TEXT_SUCCESS);
  885. }
  886. else
  887. {
  888. snprintf(data, GEARMAN_TEXT_RESPONSE_SIZE, TEXT_ERROR_CREATE_FUNCTION,
  889. (int)packet->arg_size[2], (char *)(packet->arg[2]));
  890. }
  891. }
  892. else
  893. {
  894. // create
  895. snprintf(data, GEARMAN_TEXT_RESPONSE_SIZE, TEXT_ERROR_ARGS, (int)packet->arg_size[0], (char *)(packet->arg[0]));
  896. }
  897. }
  898. else if (! strcasecmp("drop", (char *)(packet->arg[0])))
  899. {
  900. if (packet->argc == 3 && !strcasecmp("function", (char *)(packet->arg[1])))
  901. {
  902. bool success= false;
  903. for (gearman_server_function_st *function= Server->function_list; function != NULL; function= function->next)
  904. {
  905. if (! strcasecmp(function->function_name, (char *)(packet->arg[2])))
  906. {
  907. success++;
  908. if (function->worker_count == 0 && function->job_running == 0)
  909. {
  910. gearman_server_function_free(Server, function);
  911. snprintf(data, GEARMAN_TEXT_RESPONSE_SIZE, TEXT_SUCCESS);
  912. }
  913. else
  914. {
  915. snprintf(data, GEARMAN_TEXT_RESPONSE_SIZE, "ERR there are still connected workers or executing clients\r\n");
  916. }
  917. break;
  918. }
  919. }
  920. if (! success)
  921. {
  922. snprintf(data, GEARMAN_TEXT_RESPONSE_SIZE, "ERR function not found\r\n");
  923. gearmand_debug(data);
  924. }
  925. }
  926. else
  927. {
  928. // drop
  929. snprintf(data, GEARMAN_TEXT_RESPONSE_SIZE, TEXT_ERROR_ARGS, (int)packet->arg_size[0], (char *)(packet->arg[0]));
  930. }
  931. }
  932. else if (!strcasecmp("maxqueue", (char *)(packet->arg[0])))
  933. {
  934. if (packet->argc == 1)
  935. {
  936. snprintf(data, GEARMAN_TEXT_RESPONSE_SIZE, TEXT_ERROR_ARGS, (int)packet->arg_size[0], (char *)(packet->arg[0]));
  937. }
  938. else
  939. {
  940. if (packet->argc == 2)
  941. max_queue_size= GEARMAN_DEFAULT_MAX_QUEUE_SIZE;
  942. else
  943. {
  944. max_queue_size= atoi((char *)(packet->arg[2]));
  945. if (max_queue_size < 0)
  946. max_queue_size= 0;
  947. }
  948. gearman_server_function_st *function;
  949. for (function= Server->function_list;
  950. function != NULL; function= function->next)
  951. {
  952. if (strlen((char *)(packet->arg[1])) == function->function_name_size &&
  953. !memcmp(packet->arg[1], function->function_name,
  954. function->function_name_size))
  955. {
  956. function->max_queue_size= (uint32_t)max_queue_size;
  957. }
  958. }
  959. snprintf(data, GEARMAN_TEXT_RESPONSE_SIZE, TEXT_SUCCESS);
  960. }
  961. }
  962. else if (strcasecmp("getpid", (char *)(packet->arg[0])) == 0)
  963. {
  964. snprintf(data, GEARMAN_TEXT_RESPONSE_SIZE, "OK %d\n", (int)getpid());
  965. }
  966. else if (!strcasecmp("shutdown", (char *)(packet->arg[0])))
  967. {
  968. if (packet->argc == 1)
  969. {
  970. Server->shutdown= true;
  971. snprintf(data, GEARMAN_TEXT_RESPONSE_SIZE, TEXT_SUCCESS);
  972. }
  973. else if (packet->argc == 2 &&
  974. !strcasecmp("graceful", (char *)(packet->arg[1])))
  975. {
  976. Server->shutdown_graceful= true;
  977. snprintf(data, GEARMAN_TEXT_RESPONSE_SIZE, TEXT_SUCCESS);
  978. }
  979. else
  980. {
  981. // shutdown
  982. snprintf(data, GEARMAN_TEXT_RESPONSE_SIZE, TEXT_ERROR_ARGS, (int)packet->arg_size[0], (char *)(packet->arg[0]));
  983. }
  984. }
  985. else if (!strcasecmp("verbose", (char *)(packet->arg[0])))
  986. {
  987. snprintf(data, GEARMAN_TEXT_RESPONSE_SIZE, "OK %s\n", gearmand_verbose_name(Gearmand()->verbose));
  988. }
  989. else if (!strcasecmp("version", (char *)(packet->arg[0])))
  990. {
  991. snprintf(data, GEARMAN_TEXT_RESPONSE_SIZE, "OK %s\n", PACKAGE_VERSION);
  992. }
  993. else
  994. {
  995. gearmand_log_debug(GEARMAN_DEFAULT_LOG_PARAM, "Failed to find command %.*s(%llu)", packet->arg_size[0], packet->arg[0],
  996. (unsigned long long)packet->arg_size[0]);
  997. snprintf(data, GEARMAN_TEXT_RESPONSE_SIZE, TEXT_ERROR_UNKNOWN_COMMAND, (int)packet->arg_size[0], (char *)(packet->arg[0]));
  998. }
  999. server_packet= gearman_server_packet_create(server_con->thread, false);
  1000. if (server_packet == NULL)
  1001. {
  1002. gearmand_debug("free");
  1003. free(data);
  1004. return GEARMAN_MEMORY_ALLOCATION_FAILURE;
  1005. }
  1006. gearmand_packet_init(&(server_packet->packet), GEARMAN_MAGIC_TEXT, GEARMAN_COMMAND_TEXT);
  1007. server_packet->packet.magic= GEARMAN_MAGIC_TEXT;
  1008. server_packet->packet.command= GEARMAN_COMMAND_TEXT;
  1009. server_packet->packet.options.complete= true;
  1010. server_packet->packet.options.free_data= true;
  1011. if (data[0] == 0)
  1012. {
  1013. snprintf(data, GEARMAN_TEXT_RESPONSE_SIZE, TEXT_ERROR_INTERNAL_ERROR);
  1014. }
  1015. server_packet->packet.data= data;
  1016. server_packet->packet.data_size= strlen(data);
  1017. int error;
  1018. if (! (error= pthread_mutex_lock(&server_con->thread->lock)))
  1019. {
  1020. GEARMAN_FIFO_ADD(server_con->io_packet, server_packet,);
  1021. (void) pthread_mutex_unlock(&server_con->thread->lock);
  1022. }
  1023. else
  1024. {
  1025. errno= error;
  1026. gearmand_perror("pthread_mutex_lock");
  1027. assert(!"pthread_mutex_lock");
  1028. }
  1029. gearman_server_con_io_add(server_con);
  1030. return GEARMAN_SUCCESS;
  1031. }
  1032. static gearmand_error_t
  1033. _server_queue_work_data(gearman_server_job_st *server_job,
  1034. gearmand_packet_st *packet, gearman_command_t command)
  1035. {
  1036. gearman_server_client_st *server_client;
  1037. uint8_t *data;
  1038. gearmand_error_t ret;
  1039. for (server_client= server_job->client_list; server_client;
  1040. server_client= server_client->job_next)
  1041. {
  1042. if (command == GEARMAN_COMMAND_WORK_EXCEPTION && !(server_client->con->is_exceptions))
  1043. {
  1044. gearmand_debug("Dropping GEARMAN_COMMAND_WORK_EXCEPTION");
  1045. continue;
  1046. }
  1047. else if (command == GEARMAN_COMMAND_WORK_EXCEPTION)
  1048. {
  1049. gearmand_debug("Sending GEARMAN_COMMAND_WORK_EXCEPTION");
  1050. }
  1051. if (packet->data_size > 0)
  1052. {
  1053. if (packet->options.free_data &&
  1054. server_client->job_next == NULL)
  1055. {
  1056. data= (uint8_t *)(packet->data);
  1057. packet->options.free_data= false;
  1058. }
  1059. else
  1060. {
  1061. data= (uint8_t *)malloc(packet->data_size);
  1062. if (data == NULL)
  1063. {
  1064. gearmand_perror("malloc");
  1065. return GEARMAN_MEMORY_ALLOCATION_FAILURE;
  1066. }
  1067. memcpy(data, packet->data, packet->data_size);
  1068. }
  1069. }
  1070. else
  1071. {
  1072. data= NULL;
  1073. }
  1074. ret= gearman_server_io_packet_add(server_client->con, true,
  1075. GEARMAN_MAGIC_RESPONSE, command,
  1076. packet->arg[0], packet->arg_size[0],
  1077. data, packet->data_size, NULL);
  1078. if (ret != GEARMAN_SUCCESS)
  1079. {
  1080. gearmand_gerror("gearman_server_io_packet_add", ret);
  1081. return ret;
  1082. }
  1083. }
  1084. return GEARMAN_SUCCESS;
  1085. }