Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
37 changes: 37 additions & 0 deletions manpages/mxqsub.1
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,43 @@ is used to queue a job to be executed on a cluster node\&. <command> [arguments]
.RS 4
specify the estimated runtime\&.
.RE
.PP
\fB\-c \fR\fB\fI<script>\fR\fR, \fB\-\-callback=\fR\fB\fI<script>\fR\fR
.RS 4
execute \fI<script>\fR after the job reaches any terminal state (finished, failed, killed, or unknown)\&.
The script must be given as an absolute path\&.
It is run as the submitting user in the job working directory\&.
The following environment variables are set for the script:
.sp
.RS 4
\fBMXQ_JOB_ID\fR \- the job id
.br
\fBMXQ_GROUP_ID\fR \- the job group id
.br
\fBMXQ_JOB_STATUS\fR \- one of \fIfinished\fR, \fIfailed\fR, \fIkilled\fR, or \fIunknown\fR
.br
\fBMXQ_JOB_WORKDIR\fR \- the job working directory
.RE
.RE
.SH "EXAMPLES"
.PP
Send an email when a job finishes\&. Save the following script as e\&.g\&.
\fI~/bin/mxq\-notify\fR and make it executable:
.sp
.RS 4
.nf
#!/bin/sh
mail \-s "MXQ job $MXQ_JOB_ID $MXQ_JOB_STATUS" "$USER"
.fi
.RE
.sp
Then submit a job with:
.sp
.RS 4
.nf
mxqsub \-\-callback=$HOME/bin/mxq\-notify myjob
.fi
.RE
.SH "ENVIRONMENT"
.PP
\fBMXQ_MYSQL_DEFAULTFILE\fR
Expand Down
7 changes: 6 additions & 1 deletion mxq_job.c
Original file line number Diff line number Diff line change
Expand Up @@ -18,7 +18,7 @@
#include "mxq_group.h"
#include "mxq_job.h"

#define JOB_FIELDS_CNT 37
#define JOB_FIELDS_CNT 38
#define JOB_FIELDS \
" job_id, " \
" job_status, " \
Expand All @@ -29,6 +29,7 @@
" job_argv, " \
" job_stdout, " \
" job_stderr, " \
" job_callback, " \
" job_umask, " \
" host_submit, " \
" host_id, " \
Expand Down Expand Up @@ -73,6 +74,7 @@ static void bind_result_job_fields(struct mx_mysql_bind *result, struct mxq_job
mx_mysql_bind_var(result, idx++, string, &(j->job_argv_str));
mx_mysql_bind_var(result, idx++, string, &(j->job_stdout));
mx_mysql_bind_var(result, idx++, string, &(j->job_stderr));
mx_mysql_bind_var(result, idx++, string, &(j->job_callback));
mx_mysql_bind_var(result, idx++, uint32, &(j->job_umask));
mx_mysql_bind_var(result, idx++, string, &(j->host_submit));
mx_mysql_bind_var(result, idx++, string, &(j->host_id));
Expand Down Expand Up @@ -136,6 +138,7 @@ void mxq_job_free_content(struct mxq_job *j)
mx_free_null(j->job_argv_str);
mx_free_null(j->job_stdout);
mx_free_null(j->job_stderr);
mx_free_null(j->job_callback);

if (j->tmp_stderr == j->tmp_stdout) {
j->tmp_stdout = NULL;
Expand Down Expand Up @@ -605,6 +608,8 @@ int mxq_set_job_status_unknown(struct mx_mysql *mysql, struct mxq_job *job)
return res;
}

job->job_status = MXQ_JOB_STATUS_UNKNOWN;

return res;
}

Expand Down
1 change: 1 addition & 0 deletions mxq_job.h
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,7 @@ struct mxq_job {

char * job_stdout;
char * job_stderr;
char * job_callback;

char * tmp_stdout;
char * tmp_stderr;
Expand Down
50 changes: 50 additions & 0 deletions mxqd.c
Original file line number Diff line number Diff line change
Expand Up @@ -1928,6 +1928,52 @@ static void release_gpu(struct mxq_server *server, struct mxq_group *group, stru
}
}

static void run_job_callback(struct mxq_group *group, struct mxq_job *job)
{
if (!job->job_callback || !*job->job_callback)
return;

mx_log_info("job=%s(%d):%lu:%lu :: running callback: %s",
group->user_name, group->user_uid, group->group_id, job->job_id,
job->job_callback);

pid_t pid = fork();
if (pid < 0) {
mx_log_err("job=%s(%d):%lu:%lu callback fork(): %m",
group->user_name, group->user_uid, group->group_id, job->job_id);
return;
}

if (pid == 0) {
if (initgroups(group->user_name, group->user_gid) == -1)
_exit(1);
if (setregid(group->user_gid, group->user_gid) == -1)
_exit(1);
if (setreuid(group->user_uid, group->user_uid) == -1)
_exit(1);
if (chdir(job->job_workdir) == -1)
_exit(1);

mx_setenvf_forever("MXQ_JOB_ID", "%lu", job->job_id);
mx_setenvf_forever("MXQ_GROUP_ID", "%lu", job->group_id);
mx_setenv_forever("MXQ_JOB_STATUS", mxq_job_status_to_name(job->job_status));
mx_setenv_forever("MXQ_JOB_WORKDIR", job->job_workdir);

execl(job->job_callback, job->job_callback, NULL);
_exit(1);
}

int status;
if (waitpid(pid, &status, 0) == -1) {
Copy link
Contributor

@donald donald May 27, 2026

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Are we blocking the daemon here while the user script does its sleep infinity or whatever?

mx_log_err("job=%s(%d):%lu:%lu callback waitpid(): %m",
group->user_name, group->user_uid, group->group_id, job->job_id);
return;
}
if (!WIFEXITED(status) || WEXITSTATUS(status) != 0)
mx_log_warning("job=%s(%d):%lu:%lu :: callback exited with non-zero status",
group->user_name, group->user_uid, group->group_id, job->job_id);
}

static int job_has_finished(struct mxq_server *server, struct mxq_group *group, struct mxq_job_list *jlist)
{
int cnt;
Expand All @@ -1942,6 +1988,8 @@ static int job_has_finished(struct mxq_server *server, struct mxq_group *group,

rename_outfiles(server, group, job);

run_job_callback(group, job);

cnt = jlist->group->slots_per_job;
cpuset_clear_running(&server->cpu_set_running, &job->host_cpu_set);
release_gpu(server, group, job);
Expand All @@ -1966,6 +2014,8 @@ static int job_is_lost(struct mxq_server *server,struct mxq_group *group, struct

rename_outfiles(server, group, job);

run_job_callback(group, job);

cnt = jlist->group->slots_per_job;
cpuset_clear_running(&server->cpu_set_running, &job->host_cpu_set);
release_gpu(server, group, job);
Expand Down
26 changes: 24 additions & 2 deletions mxqsub.c
Original file line number Diff line number Diff line change
Expand Up @@ -60,6 +60,10 @@ static void print_usage(void)
" -e, --stderr=FILE set file to capture stderr (default: <stdout>)\n"
" -u, --umask=MASK set mode to use as umask (default: current umask)\n"
" -p, --priority=PRIORITY set priority (default: 127)\n"
" -c, --callback=SCRIPT execute SCRIPT when the job finishes (any outcome)\n"
Copy link
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Technically it doesn't need to be a SCRIPT, because we execl() whatever executable it is, right?

" runs as the submitting user in the job workdir;\n"
" MXQ_JOB_ID, MXQ_GROUP_ID, MXQ_JOB_STATUS, and\n"
" MXQ_JOB_WORKDIR are set in the environment\n"
"\n"
"Job resource information:\n"
" Scheduling is done based on the resources a job needs and\n"
Expand Down Expand Up @@ -517,6 +521,7 @@ static int add_job(struct mx_mysql *mysql, struct mxq_job *j)

" job_stdout = ?,"
" job_stderr = ?,"
" job_callback = ?,"

" job_umask = ?,"

Expand All @@ -535,8 +540,9 @@ static int add_job(struct mx_mysql *mysql, struct mxq_job *j)
mx_mysql_statement_param_bind(stmt, 4, string, &(j->job_argv_str));
mx_mysql_statement_param_bind(stmt, 5, string, &(j->job_stdout));
mx_mysql_statement_param_bind(stmt, 6, string, &(j->job_stderr));
mx_mysql_statement_param_bind(stmt, 7, uint32, &(j->job_umask));
mx_mysql_statement_param_bind(stmt, 8, string, &(j->host_submit));
mx_mysql_statement_param_bind(stmt, 7, string, &(j->job_callback));
mx_mysql_statement_param_bind(stmt, 8, uint32, &(j->job_umask));
mx_mysql_statement_param_bind(stmt, 9, string, &(j->host_submit));

res = mx_mysql_statement_execute(stmt, &num_rows);
if (res < 0) {
Expand Down Expand Up @@ -705,6 +711,7 @@ int main(int argc, char *argv[])
char arg_debug;
u_int32_t arg_tmpdir;
u_int16_t arg_gpu;
char *arg_callback;

_mx_cleanup_free_ char *current_workdir = NULL;
_mx_cleanup_free_ char *arg_stdout_absolute = NULL;
Expand Down Expand Up @@ -769,6 +776,7 @@ int main(int argc, char *argv[])
MX_OPTION_REQUIRED_ARG("prerequisites", 10),
MX_OPTION_REQUIRED_ARG("tags", 11),
MX_OPTION_NO_ARG("gpu", 12),
MX_OPTION_REQUIRED_ARG("callback", 'c'),
MX_OPTION_END
};

Expand Down Expand Up @@ -798,6 +806,7 @@ int main(int argc, char *argv[])
arg_prerequisites = "";
arg_tags = NULL;
arg_gpu = 0;
arg_callback = "";

arg_mysql_default_group = getenv("MXQ_MYSQL_DEFAULT_GROUP");
if (!arg_mysql_default_group)
Expand Down Expand Up @@ -1021,6 +1030,18 @@ int main(int argc, char *argv[])
case 12:
arg_gpu = 1;
break;

case 'c':
if (!(*optctl.optarg)) {
mx_log_crit("--callback '%s': String is empty.", optctl.optarg);
exit(EX_CONFIG);
}
if (optctl.optarg[0] != '/') {
mx_log_crit("--callback '%s': must be an absolute path.", optctl.optarg);
exit(EX_CONFIG);
}
arg_callback = optctl.optarg;
break;
}
}

Expand Down Expand Up @@ -1133,6 +1154,7 @@ int main(int argc, char *argv[])
job.job_argc = argc;
job.job_argv = argv;
job.job_argv_str = arg_args;
job.job_callback = arg_callback;

/******************************************************************/

Expand Down
1 change: 1 addition & 0 deletions mysql/create_tables.sql
Original file line number Diff line number Diff line change
Expand Up @@ -83,6 +83,7 @@ CREATE TABLE IF NOT EXISTS mxq_job (

job_stdout VARCHAR(4096) NOT NULL DEFAULT '/dev/null',
job_stderr VARCHAR(4096) NOT NULL DEFAULT '/dev/null',
job_callback VARCHAR(4096) NOT NULL DEFAULT '',

job_umask INT4 NOT NULL,

Expand Down
2 changes: 2 additions & 0 deletions mysql/migrate_019_add_job_callback.sql
Original file line number Diff line number Diff line change
@@ -0,0 +1,2 @@
ALTER TABLE mxq_job
ADD COLUMN job_callback VARCHAR(4096) NOT NULL DEFAULT '' AFTER job_stderr;