task.py 26.7 KB
Newer Older
1
# coding: utf8
2
import json
3
import os
4
import defs
5
import re
6
import os.path
7
8
9
import time
import sys
import datetime
10
11
import random
import xmlrpclib
12
import tools_utils
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28

def assert_scheduler_task_does_not_exist(args):
    ##check scheduled run
    row = db( ( db.scheduler_task.args == args)
         & ( db.scheduler_task.status != "FAILED"  )
         & ( db.scheduler_task.status != "EXPIRED"  )
         & ( db.scheduler_task.status != "TIMEOUT"  )
         & ( db.scheduler_task.status != "COMPLETED"  )
         ).select()

    if len(row) > 0 :
        res = {"message": "task already registered"}
        return res

    return None

29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
def compute_contamination(sequence_file_id, results_file_id, config_id):
    result = []
    for i in range(0, len(sequence_file_id)) :
        result.append({})
        result[i]["total_clones"]=0
        result[i]["total_reads"]=0
        result[i]["sample"]={}
        
        #open sample
        with open(defs.DIR_RESULTS+db.results_file[results_file_id[i]].data_file, "r") as json_data:
            d = json.load(json_data)
            list1 = {}
            total_reads1=d["reads"]["segmented"][0]
            for clone in d["clones"]:
                list1[clone["id"]] = clone["reads"][0]
            json_data.close()
        
        #iterate trough others run's samples
        sample_set_run = db( (db.sample_set_membership.sequence_file_id == sequence_file_id[i])
                           & (db.sample_set_membership.sample_set_id == db.sample_set.id)
                           & (db.sample_set.sample_type == "run") ).select().first()
        
        if sample_set_run != None :
            sample_set_id = sample_set_run.sample_set.id
            query = db(  ( db.sample_set.id == sample_set_id )
                       & ( db.sample_set.id == db.sample_set_membership.sample_set_id )
                       & ( db.sequence_file.id == db.sample_set_membership.sequence_file_id)
                       & ( db.sequence_file.id != sequence_file_id[i])
                       & ( db.results_file.sequence_file_id == db.sequence_file.id )
                       & ( db.results_file.config_id == config_id[i]  )
                       ).select(db.sequence_file.ALL,db.results_file.ALL, db.sample_set.id, orderby=db.sequence_file.id|~db.results_file.run_date)

            query2 = []
            sfi = 0
            for row in query : 
                if row.sequence_file.id != sfi :
                    query2.append(row)
                    sfi = row.sequence_file.id
            
            for row in query2 :
                print defs.DIR_RESULTS+row.results_file.data_file
                result[i]["sample"][row.results_file.id] = {}
                result[i]["sample"][row.results_file.id]["clones"] = 0
                result[i]["sample"][row.results_file.id]["reads"] = 0
                with open(defs.DIR_RESULTS+row.results_file.data_file, "r") as json_data2:
                    try:
                        d = json.load(json_data2)
                        total_reads2=d["reads"]["segmented"][0]
                        for clone in d["clones"]:
                            if clone["id"] in list1 :
                                if clone["reads"][0] > 10*list1[clone["id"]] :
                                    result[i]["total_clones"] += 1
                                    result[i]["total_reads"] += list1[clone["id"]]
                                    result[i]["sample"][row.results_file.id]["clones"] += 1
                                    result[i]["sample"][row.results_file.id]["reads"] += list1[clone["id"]]
                            
                    except ValueError, e:
                        print 'invalid_json'
                    json_data2.close()


    return result
    
92

93
def schedule_run(id_sequence, id_config, grep_reads=None):
94
    from subprocess import Popen, PIPE, STDOUT, os
95

96
97
98
    #check results_file
    row = db( ( db.results_file.config_id == id_config ) & 
             ( db.results_file.sequence_file_id == id_sequence )  
99
100
             ).select()
    
Marc Duez's avatar
Marc Duez committed
101
102
103
104
    ts = time.time()
    data_id = db.results_file.insert(sequence_file_id = id_sequence,
                                     run_date = datetime.datetime.fromtimestamp(ts).strftime('%Y-%m-%d %H:%M:%S'),
                                     config_id = id_config )
105
        
106
    args = [id_sequence, id_config, data_id, grep_reads]
107

108
109
110
111
    err = assert_scheduler_task_does_not_exist(str(args))
    if err:
        log.error(err)
        return err
112

113
    program = db.config[id_config].program
114
    ##add task to scheduler
115
116
    task = scheduler.queue_task(program, args,
                                repeats = 1, timeout = defs.TASK_TIMEOUT)
117
118
119
120
    
    if db.sequence_file[id_sequence].pre_process_flag == "WAIT" or db.sequence_file[id_sequence].pre_process_flag == "RUN" : 
        db.scheduler_task[task.id] = dict(status ="STOPPED")
    
121
    db.results_file[data_id] = dict(scheduler_task_id = task.id)
122

123
    filename= db.sequence_file[id_sequence].filename
124

125
    res = {"redirect": "reload",
126
           "message": "[%s] c%s: process requested - %s %s" % (data_id, id_config, grep_reads, filename)}
127

128
    log.info(res)
129
130
    return res

131

132
def run_vidjil(id_file, id_config, id_data, grep_reads,
133
               clean_before=False, clean_after=False):
134
    from subprocess import Popen, PIPE, STDOUT, os
135
136
137
    from datetime import timedelta as timed
    
    ## re schedule if pre_process is still pending
marc's avatar
marc committed
138
    if db.sequence_file[id_file].pre_process_flag == "WAIT" or db.sequence_file[id_file].pre_process_flag == "RUN" :
139
        
140
        print "Pre-process is still pending, re-schedule"
141
    
142
        args = [id_file, id_config, id_data, grep_reads]
143
144
        task = scheduler.queue_task("vidjil", args,
                        repeats = 1, timeout = defs.TASK_TIMEOUT,
145
                               start_time=request.now + timed(seconds=1200))
146
147
148
149
150
151
152
        db.results_file[id_data] = dict(scheduler_task_id = task.id)
        db.commit()
        print task.id
        sys.stdout.flush()
        
        return "SUCCESS"
    
marc's avatar
marc committed
153
154
155
156
    if db.sequence_file[id_file].pre_process_flag == "FAILED" :
        print "Pre-process has failed"
        raise ValueError('pre-process has failed')
        return "FAIL"
157
158
    
    ## les chemins d'acces a vidjil / aux fichiers de sequences
159
    germline_folder = defs.DIR_GERMLINE
160
161
    upload_folder = defs.DIR_SEQUENCES
    out_folder = defs.DIR_OUT_VIDJIL_ID % id_data
162
    
163
164
165
166
    cmd = "rm -rf "+out_folder 
    p = Popen(cmd, shell=True, stdin=PIPE, stdout=PIPE, stderr=STDOUT, close_fds=True)
    p.wait()
    
167
168
169
    ## filepath du fichier de séquence
    row = db(db.sequence_file.id==id_file).select()
    filename = row[0].data_file
170
    output_filename = defs.BASENAME_OUT_VIDJIL_ID % id_data
171
    seq_file = upload_folder+filename
172

173
    ## config de vidjil
174
    vidjil_cmd = db.config[id_config].command
175
    vidjil_cmd = vidjil_cmd.replace( ' germline' ,germline_folder)
176
177
178
179

    if grep_reads:
        # TODO: security, assert grep_reads XXXX
        vidjil_cmd += ' -FaW "%s" ' % grep_reads
180
    
181
    os.makedirs(out_folder)
182
183
    out_log = out_folder+'/'+output_filename+'.vidjil.log'
    vidjil_log_file = open(out_log, 'w')
184

185
    ## commande complete
186
    cmd = defs.DIR_VIDJIL + '/vidjil ' + ' -o  ' + out_folder + " -b " + output_filename
187
    cmd += ' ' + vidjil_cmd + ' '+ seq_file
188
    
189
    ## execute la commande vidjil
190
191
192
193
194
    print "=== Launching Vidjil ==="
    print cmd    
    print "========================"
    sys.stdout.flush()

195
    p = Popen(cmd, shell=True, stdin=PIPE, stdout=vidjil_log_file, stderr=STDOUT, close_fds=True)
196

197
    (stdoutdata, stderrdata) = p.communicate()
198

199
    print "Output log in " + out_log
200
201
202
    sys.stdout.flush()
    db.commit()

203
204
205
206
207
208
209
210
    ## Get result file
    if grep_reads:
        out_results = out_folder + '/seq/clone.fa-1'
    else:
        out_results = out_folder + '/' + output_filename + '.vidjil'

    print "===>", out_results
    results_filepath = os.path.abspath(out_results)
211
212
213
214

    try:
        stream = open(results_filepath, 'rb')
    except IOError:
215
        print "!!! Vidjil failed, no result file"
216
217
        res = {"message": "[%s] c%s: Vidjil FAILED - %s" % (id_data, id_config, out_folder)}
        log.error(res)
218
        raise IOError
219
    
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
    ## Parse some info in .log
    vidjil_log_file.close()

    segmented = re.compile("==> segmented (\d+) reads \((\d*\.\d+|\d+)%\)")
    windows = re.compile("==> found (\d+) .*-windows in .* segments .* inside (\d+) sequences")
    info = ''
    reads = None
    segs = None
    ratio = None
    wins = None
    for l in open(out_log):
        m = segmented.search(l)
        if m:
            print l,
            segs = int(m.group(1))
            ratio = m.group(2)
            info = "%d segmented (%s%%)" % (segs, ratio)
            continue
        m = windows.search(l)
        if m:
            print l,
            wins = int(m.group(1))
            reads = int(m.group(2))
            info = "%d reads, " % reads + info + ", %d windows" % wins
            break


247
    ## insertion dans la base de donnée
248
249
    ts = time.time()
    
250
    db.results_file[id_data] = dict(status = "ready",
251
                                 run_date = datetime.datetime.fromtimestamp(ts).strftime('%Y-%m-%d %H:%M:%S'),
252
253
                                 data_file = stream
                                )
254
    
255
256
    db.commit()
    
Mathieu Giraud's avatar
Mathieu Giraud committed
257
258
259
260
    if clean_after:
        clean_cmd = "rm -rf " + out_folder 
        p = Popen(clean_cmd, shell=True, stdin=PIPE, stdout=PIPE, stderr=STDOUT, close_fds=True)
        p.wait()
261
    
Mathieu Giraud's avatar
Mathieu Giraud committed
262
263
    ## l'output de Vidjil est stocké comme resultat pour l'ordonnanceur
    ## TODO parse result success/fail
264

265
266
    config_name = db.config[id_config].name

267
    res = {"message": "[%s] c%s: Vidjil finished - %s - %s" % (id_data, id_config, info, out_folder)}
268
269
    log.info(res)

270
271

    if not grep_reads:
272
273
274
        for row in db(db.sample_set_membership.sequence_file_id==id_file).select() :
	    sample_set_id = row.sample_set_id
	    print row.sample_set_id
275
            run_fuse(id_file, id_config, id_data, sample_set_id, clean_before = False)
276

Mathieu Giraud's avatar
Mathieu Giraud committed
277
    return "SUCCESS"
278

279
def run_mixcr(id_file, id_config, id_data, clean_before=False, clean_after=False):
280
281
    from subprocess import Popen, PIPE, STDOUT, os
    import time
282
    import json
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301

    upload_folder = defs.DIR_SEQUENCES
    out_folder = defs.DIR_OUT_VIDJIL_ID % id_data

    cmd = "rm -rf "+out_folder
    p = Popen(cmd, shell=True, stdin=PIPE, stdout=PIPE, stderr=STDOUT, close_fds=True)
    p.wait()

    ## filepath du fichier de séquence
    row = db(db.sequence_file.id==id_file).select()
    filename = row[0].data_file
    output_filename = defs.BASENAME_OUT_VIDJIL_ID % id_data
    seq_file = upload_folder+filename

    ## config de vidjil
    arg_cmd = db.config[id_config].command

    os.makedirs(out_folder)
    out_log = out_folder + output_filename+'.mixcr.log'
302
    report = out_folder + output_filename + '.mixcr.report'
303
304
305
306
    log_file = open(out_log, 'w')

    out_alignments = out_folder + output_filename + '.align.vdjca'
    out_clones =  out_folder + output_filename + '.clones.clns'
307
308
    out_results_file = output_filename + '.mixcr'
    out_results = out_folder + out_results_file
309

310
    align_report = report + '.aln'
311
    assembly_report = report + '.asmbl'
312

313
314
315
316
317
318
319
320
    arg_cmds = arg_cmd.split('|')
    args_1, args_2, args_3 = '', '', ''
    try:
        args_1, args_2, args_3 = arg_cmds
    except:
        print arg_cmd
        print "! Bad arguments, we expect args_align | args_assemble | args_exportClones"

321
322
    ## commande complete
    mixcr = defs.DIR_MIXCR + 'mixcr'
323
324
325
    cmd = mixcr + ' align --save-reads -t 1 -r ' + align_report + ' ' + args_1 + ' ' + seq_file  + ' ' + out_alignments
    cmd += ' && '
    cmd += mixcr + ' assemble -t 1 -r ' + assembly_report + ' ' + args_2 + ' ' + out_alignments + ' ' + out_clones
326
    cmd += ' && rm ' + out_alignments
327
328
    cmd += ' && '
    cmd += mixcr + ' exportClones --format vidjil -germline -id -name -reads -sequence -top -seg -s ' + args_3 + ' ' + out_clones + ' ' + out_results
329
330
331
332
333
334
335
336

    ## execute la commande MiXCR
    print "=== Launching MiXCR ==="
    print cmd
    print "========================"
    sys.stdout.flush()

    p = Popen(cmd, shell=True, stdin=PIPE, stdout=log_file, stderr=STDOUT, close_fds=True)
337
    p.wait()
338
339
340
341
342
343

    print "Output log in " + out_log
    sys.stdout.flush()

    ## Get result file
    print "===>", out_results
HERBERT Ryan's avatar
HERBERT Ryan committed
344
345
346
    results_filepath = os.path.abspath(out_results)
    try:
        stream = open(results_filepath, 'rb')
347
        stream.close()
HERBERT Ryan's avatar
HERBERT Ryan committed
348
349
350
351
352
    except IOError:
        print "!!! MiXCR failed, no result file"
        res = {"message": "[%s] c%s: MiXCR FAILED - %s" % (id_data, id_config, out_folder)}
        log.error(res)
        raise IOError
353

HERBERT Ryan's avatar
HERBERT Ryan committed
354
355
356
    align_report = get_file_content(align_report)
    assembly_report = get_file_content(assembly_report)
    reports = align_report + assembly_report
357
    original_name = row[0].filename
HERBERT Ryan's avatar
HERBERT Ryan committed
358
    totalReads = extract_total_reads(assembly_report)
359
360
    with open(results_filepath, 'r') as json_file:
        my_json = json.load(json_file)
HERBERT Ryan's avatar
HERBERT Ryan committed
361
362
        fill_field(my_json, reports, "log", "samples", True)
        fill_field(my_json, original_name, "original_names", "samples")
363
        fill_field(my_json, cmd, "commandline", "samples")
HERBERT Ryan's avatar
HERBERT Ryan committed
364
        fill_field(my_json, totalReads, "total", "reads")
365

366
367
368
369
        # TODO fix this dirty hack to get around bad file descriptor error
    new_file = open(results_filepath, 'w')
    json.dump(my_json, new_file)
    new_file.close()
370

371
372
    ## insertion dans la base de donnée
    ts = time.time()
HERBERT Ryan's avatar
HERBERT Ryan committed
373
    
374
    stream = open(results_filepath, 'rb')
375
376
    db.results_file[id_data] = dict(status = "ready",
                                 run_date = datetime.datetime.fromtimestamp(ts).strftime('%Y-%m-%d %H:%M:%S'),
HERBERT Ryan's avatar
HERBERT Ryan committed
377
                                 data_file = stream
378
                                )
HERBERT Ryan's avatar
HERBERT Ryan committed
379
    
380
381
382
383
384
385
    db.commit()

    config_name = db.config[id_config].name
    res = {"message": "[%s] c%s: MiXCR finished - %s" % (id_data, id_config, out_folder)}
    log.info(res)

HERBERT Ryan's avatar
HERBERT Ryan committed
386
387
388
    for row in db(db.sample_set_membership.sequence_file_id==id_file).select() :
        sample_set_id = row.sample_set_id
        print row.sample_set_id
389
        run_fuse(id_file, id_config, id_data, sample_set_id, clean_before = False)
HERBERT Ryan's avatar
HERBERT Ryan committed
390
391


392
    return "SUCCESS"
393

394
def run_copy(id_file, id_config, id_data, clean_before=False, clean_after=False):
395
396
397
398
    from subprocess import Popen, PIPE, STDOUT, os
    
    ## les chemins d'acces a vidjil / aux fichiers de sequences
    upload_folder = defs.DIR_SEQUENCES
399
    output_filename = defs.BASENAME_OUT_VIDJIL_ID % id_data
400
401
402
403
404
405
406
407
408
409
410
411
    out_folder = defs.DIR_OUT_VIDJIL_ID % id_data
    
    cmd = "rm -rf "+out_folder 
    p = Popen(cmd, shell=True, stdin=PIPE, stdout=PIPE, stderr=STDOUT, close_fds=True)
    p.wait()
    
    ## filepath du fichier de séquence
    row = db(db.sequence_file.id==id_file).select()
    filename = row[0].data_file
    
    os.makedirs(out_folder)
    vidjil_log_file = open(out_folder+'/'+output_filename+'.vidjil.log', 'w')
412

413
414
415
416
417
418
    print "Output log in "+out_folder+'/'+output_filename+'.vidjil.log'
    sys.stdout.flush()
    db.commit()
    
    ## récupération du fichier 
    results_filepath = os.path.abspath(defs.DIR_SEQUENCES+row[0].data_file)
419

420
421
422
    try:
        stream = open(results_filepath, 'rb')
    except IOError:
423
424
425
        print "!!! 'copy' failed, no file"
        res = {"message": "[%s] c%s: 'copy' FAILED - %s - %s" % (id_data, id_config, info, out_folder)}
        log.error(res)
426
427
428
429
430
431
432
433
        raise IOError
    
    ## insertion dans la base de donnée
    ts = time.time()
    
    db.results_file[id_data] = dict(status = "ready",
                                 run_date = datetime.datetime.fromtimestamp(ts).strftime('%Y-%m-%d %H:%M:%S'),
                                 #data_file = row[0].data_file
434
                                 data_file = db.results_file.data_file.store(stream, row[0].filename)
435
436
437
438
439
440
441
442
443
444
445
                                )
    db.commit()
    
    if clean_after:
        clean_cmd = "rm -rf " + out_folder 
        p = Popen(clean_cmd, shell=True, stdin=PIPE, stdout=PIPE, stderr=STDOUT, close_fds=True)
        p.wait()
    
    ## l'output de Vidjil est stocké comme resultat pour l'ordonnanceur
    ## TODO parse result success/fail

446
    res = {"message": "[%s] c%s: 'copy' finished - %s" % (id_data, id_config, filename)}
447
448
    log.info(res)

449
450
451
    for row in db(db.sample_set_membership.sequence_file_id==id_file).select() :
        sample_set_id = row.sample_set_id
        print row.sample_set_id
452
        run_fuse(id_file, id_config, id_data, sample_set_id, clean_before = False)
453

454
455

    return "SUCCESS"
Mathieu Giraud's avatar
Mathieu Giraud committed
456
457
458



459
def run_fuse(id_file, id_config, id_data, sample_set_id, clean_before=True, clean_after=False):
460
461
    from subprocess import Popen, PIPE, STDOUT, os
    
462
    out_folder = defs.DIR_OUT_VIDJIL_ID % id_data
463
    output_filename = defs.BASENAME_OUT_VIDJIL_ID % id_data + '-%s' % sample_set_id
464
    
Mathieu Giraud's avatar
Mathieu Giraud committed
465
466
467
468
469
    if clean_before:
        cmd = "rm -rf "+out_folder 
        p = Popen(cmd, shell=True, stdin=PIPE, stdout=PIPE, stderr=STDOUT, close_fds=True)
        p.wait()
        os.makedirs(out_folder)    
470
471
    
    
Mathieu Giraud's avatar
Mathieu Giraud committed
472
    fuse_log_file = open(out_folder+'/'+output_filename+'.fuse.log', 'w')
Marc's avatar
Marc committed
473
    
474
    ## fuse.py 
475
    output_file = out_folder+'/'+output_filename+'.fused'
476
    files = ""
477
    sequence_file_list = ""
478
    query2 = db( ( db.results_file.sequence_file_id == db.sequence_file.id )
479
                   & ( db.sample_set_membership.sequence_file_id == db.sequence_file.id)
480
                   & ( db.sample_set_membership.sample_set_id == sample_set_id)
481
                   & ( db.results_file.config_id == id_config )
Vidjil Team's avatar
Vidjil Team committed
482
                   ).select( orderby=db.sequence_file.id|~db.results_file.run_date) 
Marc Duez's avatar
Marc Duez committed
483
484
485
486
487
488
489
    query = []
    sequence_file_id = 0
    for row in query2 : 
        if row.sequence_file.id != sequence_file_id :
            query.append(row)
            sequence_file_id = row.sequence_file.id
            
490
    for row in query :
Mathieu Giraud's avatar
Mathieu Giraud committed
491
492
        if row.results_file.data_file is not None :
            files += defs.DIR_RESULTS + row.results_file.data_file + " "
493
494
            sequence_file_list += str(row.results_file.sequence_file_id) + "_"
            
495
496
    fuse_cmd = db.config[id_config].fuse_command
    cmd = "python "+defs.DIR_FUSE+"/fuse.py -o "+ output_file + " " + fuse_cmd + " " + files
Mathieu Giraud's avatar
Mathieu Giraud committed
497
498
499


    print "=== fuse.py ==="
500
    print cmd
Mathieu Giraud's avatar
Mathieu Giraud committed
501
502
503
504
505
506
507
    print "==============="
    sys.stdout.flush()

    p = Popen(cmd, shell=True, stdin=PIPE, stdout=fuse_log_file, stderr=STDOUT, close_fds=True)
    (stdoutdata, stderrdata) = p.communicate()
    print "Output log in "+out_folder+'/'+output_filename+'.fuse.log'

508
    fuse_filepath = os.path.abspath(output_file)
509
510
511
512
513

    try:
        stream = open(fuse_filepath, 'rb')
    except IOError:
        print "!!! Fuse failed, no .fused file"
514
515
        res = {"message": "[%s] c%s: 'fuse' FAILED - %s" % (id_data, id_config, output_file)}
        log.error(res)
516
517
        raise IOError

518
    ts = time.time()
519
520
521
522
523
524
525
526
527
528

    fused_files = db( ( db.fused_file.config_id == id_config ) &
                     ( db.fused_file.sample_set_id == sample_set_id )
                 ).select()
    if len(fused_files) > 0:
        id_fuse = fused_files[0].id
    else:
        id_fuse = db.fused_file.insert(sample_set_id = sample_set_id,
                                       config_id = id_config)

529
    db.fused_file[id_fuse] = dict(fuse_date = datetime.datetime.fromtimestamp(ts).strftime('%Y-%m-%d %H:%M:%S'),
530
531
                                 fused_file = stream,
                                 sequence_file_list = sequence_file_list)
532
    db.commit()
Mathieu Giraud's avatar
Mathieu Giraud committed
533
534
535
536
537

    if clean_after:
        clean_cmd = "rm -rf " + out_folder 
        p = Popen(clean_cmd, shell=True, stdin=PIPE, stdout=PIPE, stderr=STDOUT, close_fds=True)
        p.wait()
538
    
539
    res = {"message": "[%s] c%s: 'fuse' finished - %s" % (id_data, id_config, output_file)}
540
541
    log.info(res)

542
    return "SUCCESS"
543

544
def custom_fuse(file_list):
545
    from subprocess import Popen, PIPE, STDOUT, os
546
547
548

    if defs.PORT_FUSE_SERVER is None:
        raise IOError('This server cannot fuse custom data')
549
    random_id = random.randint(99999999,99999999999)
550
    out_folder = os.path.abspath(defs.DIR_OUT_VIDJIL_ID % random_id)
551
    output_filename = defs.BASENAME_OUT_VIDJIL_ID % random_id
552
553
554
555
556
    
    cmd = "rm -rf "+out_folder 
    p = Popen(cmd, shell=True, stdin=PIPE, stdout=PIPE, stderr=STDOUT, close_fds=True)
    p.wait()
    os.makedirs(out_folder)    
557

558
    res = {"message": "'custom fuse' (%d files): %s" % (len(file_list), ','.join(file_list))}
559
560
    log.info(res)
        
561
562
563
564
565
    ## fuse.py 
    output_file = out_folder+'/'+output_filename+'.fused'
    files = ""
    for id in file_list :
        if db.results_file[id].data_file is not None :
566
            files += os.path.abspath(defs.DIR_RESULTS + db.results_file[id].data_file) + " "
567
    
568
    cmd = "python "+ os.path.abspath(defs.DIR_FUSE) +"/fuse.py -o "+output_file+" -t 100 "+files
HERBERT Ryan's avatar
HERBERT Ryan committed
569
    proc_srvr = xmlrpclib.ServerProxy("https://localhost:%d" % defs.PORT_FUSE_SERVER)
570
    fuse_filepath = proc_srvr.fuse(cmd, out_folder, output_filename)
571
572
573
574
    
    try:
        f = open(fuse_filepath, 'rb')
        data = gluon.contrib.simplejson.loads(f.read())
575
    except IOError, e:
576
577
        res = {"message": "'custom fuse' -> IOError"}
        log.error(res)
578
        raise e
579
580
581
582
583

    clean_cmd = "rm -rf " + out_folder 
    p = Popen(clean_cmd, shell=True, stdin=PIPE, stdout=PIPE, stderr=STDOUT, close_fds=True)
    p.wait()

584
585
586
    res = {"message": "'custom fuse' -> finished"}
    log.info(res)

587
    return data
588

HERBERT Ryan's avatar
HERBERT Ryan committed
589
#TODO move this ?
590
591
592
593
594
595
def get_file_content(filename):
    content = ""
    with open(filename, 'rb') as my_file:
        content = my_file.read()
    return content

HERBERT Ryan's avatar
HERBERT Ryan committed
596
597
598
599
600
601
602
603
604
605
606
607
608
def fill_field(dest, value, field, parent="", append=False):
    if parent is not None and parent != "":
        dest = dest[parent]

    if field in dest:
        if append:
            dest[field][0] += value
        else:
            dest[field][0] = value
    else:
        dest[field] = []
        dest[field].append(value)

HERBERT Ryan's avatar
HERBERT Ryan committed
609
610
611
612
def extract_total_reads(report):
    reads_matcher = re.compile("Total Reads analysed: [0-9]+")
    reads_line = reads_matcher.search(report).group()
    return reads_line.split(' ')[-1]
613
614
615
616
617
618
619
620
621
622
623
624
625
626

def schedule_pre_process(sequence_file_id, pre_process_id):
    from subprocess import Popen, PIPE, STDOUT, os


    args = [pre_process_id, sequence_file_id]

    err = assert_scheduler_task_does_not_exist(str(args))
    if err:
        log.error(err)
        return err

    task = scheduler.queue_task("pre_process", args,
                                repeats = 1, timeout = defs.TASK_TIMEOUT)
627
628
    db.sequence_file[sequence_file_id] = dict(pre_process_scheduler_task_id = task.id)
    
629
    res = {"redirect": "reload",
630
           "message": "{%s} (%s): process requested" % (sequence_file_id, pre_process_id)}
631
632
633
634
635

    log.info(res)
    return res


636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
def get_preprocessed_filename(filename1, filename2):
    '''
    Get the same extension and then get the most common among them

    >>> get_preprocessed_filename('lsdkj_toto-tata_mlkmfsdlkf.fastq.gz', 'rete_toto-tata_eekjdf.fastq.gz')
    '_toto-tata_.fastq.gz'
    '''
    if not filename2:
        return filename1
    extension = tools_utils.get_common_suffpref([filename1, filename2], min(len(filename1), len(filename2)), -1)
    without_extension = [x[0:-len(extension)] for x in [filename1, filename2]]
    common = tools_utils.common_substring(without_extension)
    if len(common) == 0:
        common = '_'.join(without_extension)
    return common+extension


653
def run_pre_process(pre_process_id, sequence_file_id, clean_before=True, clean_after=False):
Mathieu Giraud's avatar
Mathieu Giraud committed
654
655
656
657
658
    '''
    Run a pre-process on sequence_file.data_file (and possibly sequence_file.data_file+2),
    put the output back in sequence_file.data_file.
    '''

659
660
    from subprocess import Popen, PIPE, STDOUT, os
    
661
    sequence_file = db.sequence_file[sequence_file_id]
662
663
664
    db.sequence_file[sequence_file_id] = dict(pre_process_flag = "RUN")
    db.commit()
    
665
    out_folder = defs.DIR_PRE_VIDJIL_ID % sequence_file_id
666
667
    output_filename = get_preprocessed_filename(get_original_filename(sequence_file.data_file),
                                                get_original_filename(sequence_file.data_file2))
668
669
670
671
672
673
674
    
    if clean_before:
        cmd = "rm -rf "+out_folder 
        p = Popen(cmd, shell=True, stdin=PIPE, stdout=PIPE, stderr=STDOUT, close_fds=True)
        p.wait()
        os.makedirs(out_folder)    

675
    output_file = out_folder+'/'+output_filename
676
677
678
679
680
            
    pre_process = db.pre_process[pre_process_id]
    
    
    cmd = pre_process.command.replace( "&file1&", defs.DIR_SEQUENCES + sequence_file.data_file)
681
682
    if sequence_file.data_file2:
        cmd = cmd.replace( "&file2&", defs.DIR_SEQUENCES + sequence_file.data_file2)
683
684
    cmd = cmd.replace( "&result&", output_file)

685
    print "=== Pre-process %s ===" % pre_process_id
686
    print cmd
687
    print "==============="
688
    sys.stdout.flush()
689

690
691
    out_log = out_folder+'/'+output_filename+'.pre.log'
    log_file = open(out_log, 'w')
692
    
693
    os.chdir(defs.DIR_FUSE)
694
695
    p = Popen(cmd, shell=True, stdin=PIPE, stdout=log_file, stderr=log_file, close_fds=True)
    (stdoutdata, stderrdata) = p.communicate()
696
    print "Output log in " + out_log
697
698
699
700
701
702
703

    filepath = os.path.abspath(output_file)

    try:
        stream = open(filepath, 'rb')
    except IOError:
        print "!!! Pre-process failed, no result file"
704
        res = {"message": "{%s} p%s: 'pre_process' FAILED - %s" % (sequence_file_id, pre_process_id, output_file)}
705
        log.error(res)
706
707
        db.sequence_file[sequence_file_id] = dict(pre_process_flag = "FAILED")
        db.commit()
708
709
        raise IOError

marc's avatar
marc committed
710
711
        

Mathieu Giraud's avatar
Mathieu Giraud committed
712
713
    # Now we update the sequence file with the result of the pre-process
    # We forget the initial data_file (and possibly data_file2)
714
    db.sequence_file[sequence_file_id] = dict(data_file = stream,
715
                                              data_file2 = None,
716
                                              pre_process_flag = "DONE")
717
718
719
720
721
722
723
724
725
726
727
728
    
    #resume STOPPED task for this sequence file
    stopped_task = db(
        (db.results_file.sequence_file_id == sequence_file_id)
        & (db.results_file.scheduler_task_id == db.scheduler_task.id)
        & (db.scheduler_task.status == "STOPPED")
    ).select()
    
    for row in stopped_task :
        db.scheduler_task[row.scheduler_task.id] = dict(status = "QUEUED")
        
    
729
730
    db.commit()

731
    # Dump log in scheduler_run.run_output
732
733
    log_file.close()
    for l in open(out_log):
734
        print l,
735
736
737

    # Remove data file from disk to save space (it is now saved elsewhere)
    os.remove(filepath)
738
    
739
740
741
742
743
    if clean_after:
        clean_cmd = "rm -rf " + out_folder 
        p = Popen(clean_cmd, shell=True, stdin=PIPE, stdout=PIPE, stderr=STDOUT, close_fds=True)
        p.wait()
    
744
    res = {"message": "{%s} p%s: 'pre_process' finished - %s" % (sequence_file_id, pre_process_id, output_file)}
745
746
747
748
    log.info(res)

    return "SUCCESS"
    
HERBERT Ryan's avatar
HERBERT Ryan committed
749
    
750
from gluon.scheduler import Scheduler
751
scheduler = Scheduler(db, dict(vidjil=run_vidjil,
752
                               compute_contamination=compute_contamination,
753
                               mixcr=run_mixcr,
754
755
                               none=run_copy,
                               pre_process=run_pre_process))