editor.py 29 KB


  1. #!/usr/bin/env python
  2. # Licensed to Cloudera, Inc. under one
  3. # or more contributor license agreements. See the NOTICE file
  4. # distributed with this work for additional information
  5. # regarding copyright ownership. Cloudera, Inc. licenses this file
  6. # to you under the Apache License, Version 2.0 (the
  7. # "License"); you may not use this file except in compliance
  8. # with the License. You may obtain a copy of the License at
  9. #
  10. # http://www.apache.org/licenses/LICENSE-2.0
  11. #
  12. # Unless required by applicable law or agreed to in writing, software
  13. # distributed under the License is distributed on an "AS IS" BASIS,
  14. # WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
  15. # See the License for the specific language governing permissions and
  16. # limitations under the License.
  17. try:
  18. import json
  19. except ImportError:
  20. import simplejson as json
  21. import logging
  22. import shutil
  23. from django.core.urlresolvers import reverse
  24. from django.db.models import Q
  25. from django.forms.formsets import formset_factory
  26. from django.forms.models import inlineformset_factory
  27. from django.http import HttpResponse
  28. from django.shortcuts import redirect
  29. from django.utils.functional import curry
  30. from django.utils.translation import ugettext as _, activate as activate_translation
  31. from desktop.lib.django_util import render, extract_field_data
  32. from desktop.lib.exceptions_renderable import PopupException
  33. from desktop.lib.rest.http_client import RestException
  34. from hadoop.fs.exceptions import WebHdfsException
  35. from liboozie.submittion import Submission
  36. from filebrowser.lib.archives import archive_factory
  37. from oozie.conf import SHARE_JOBS
  38. from oozie.decorators import check_job_access_permission, check_job_edition_permission,\
  39. check_dataset_access_permission, check_dataset_edition_permission
  40. from oozie.import_workflow import import_workflow as _import_workflow
  41. from oozie.management.commands import oozie_setup
  42. from oozie.models import Workflow, History, Coordinator,\
  43. Dataset, DataInput, DataOutput,\
  44. ACTION_TYPES, Bundle, BundledCoordinator, Job
  45. from oozie.forms import WorkflowForm, CoordinatorForm, DatasetForm,\
  46. DataInputForm, DataOutputForm, LinkForm,\
  47. DefaultLinkForm, ParameterForm, ImportWorkflowForm,\
  48. NodeForm, BundleForm, BundledCoordinatorForm, design_form_by_type
  49. LOG = logging.getLogger(__name__)
  50. def list_workflows(request):
  51. show_setup_app = True
  52. data = Workflow.objects.filter(managed=True)
  53. if not SHARE_JOBS.get() and not request.user.is_superuser:
  54. data = data.filter(owner=request.user)
  55. else:
  56. data = data.filter(Q(is_shared=True) | Q(owner=request.user))
  57. data = data.order_by('-last_modified')
  58. return render('editor/list_workflows.mako', request, {
  59. 'jobs': list(data),
  60. 'json_jobs': json.dumps(list(data.values_list('id', flat=True))),
  61. 'show_setup_app': show_setup_app,
  62. })
  63. def list_coordinators(request, workflow_id=None):
  64. data = Coordinator.objects
  65. if workflow_id is not None:
  66. data = data.filter(workflow__id=workflow_id)
  67. if not SHARE_JOBS.get() and not request.user.is_superuser:
  68. data = data.filter(owner=request.user)
  69. else:
  70. data = data.filter(Q(is_shared=True) | Q(owner=request.user))
  71. data = data.order_by('-last_modified')
  72. return render('editor/list_coordinators.mako', request, {
  73. 'jobs': list(data),
  74. 'json_jobs': json.dumps(list(data.values_list('id', flat=True))),
  75. })
  76. def list_bundles(request):
  77. data = Bundle.objects
  78. if not SHARE_JOBS.get() and not request.user.is_superuser:
  79. data = data.filter(owner=request.user)
  80. else:
  81. data = data.filter(Q(is_shared=True) | Q(owner=request.user))
  82. data = data.order_by('-last_modified')
  83. return render('editor/list_bundles.mako', request, {
  84. 'jobs': list(data),
  85. 'json_jobs': json.dumps(list(data.values_list('id', flat=True))),
  86. })
  87. def create_workflow(request):
  88. workflow = Workflow.objects.new_workflow(request.user)
  89. if request.method == 'POST':
  90. workflow_form = WorkflowForm(request.POST, instance=workflow)
  91. if workflow_form.is_valid():
  92. wf = workflow_form.save()
  93. wf.managed = True
  94. Workflow.objects.initialize(wf, request.fs)
  95. return redirect(reverse('oozie:edit_workflow', kwargs={'workflow': workflow.id}))
  96. else:
  97. request.error(_('Errors on the form: %s') % workflow_form.errors)
  98. else:
  99. workflow_form = WorkflowForm(instance=workflow)
  100. return render('editor/create_workflow.mako', request, {
  101. 'workflow_form': workflow_form,
  102. 'workflow': workflow,
  103. })
  104. def import_workflow(request):
  105. workflow = Workflow.objects.new_workflow(request.user)
  106. if request.method == 'POST':
  107. workflow_form = ImportWorkflowForm(request.POST, request.FILES, instance=workflow)
  108. if workflow_form.is_valid():
  109. if workflow_form.cleaned_data.get('resource_archive'):
  110. # Upload resources to workspace
  111. source = workflow_form.cleaned_data.get('resource_archive')
  112. if source.name.endswith('.zip'):
  113. workflow.save()
  114. Workflow.objects.initialize(workflow, request.fs)
  115. temp_path = archive_factory(source).extract()
  116. request.fs.copyFromLocal(temp_path, workflow.deployment_dir)
  117. shutil.rmtree(temp_path)
  118. else:
  119. raise PopupException(_('Archive should be a Zip.'))
  120. workflow.managed = True
  121. workflow.save()
  122. workflow_definition = workflow_form.cleaned_data['definition_file'].read()
  123. try:
  124. _import_workflow(fs=request.fs, workflow=workflow, workflow_definition=workflow_definition)
  125. request.info(_('Workflow imported'))
  126. return redirect(reverse('oozie:edit_workflow', kwargs={'workflow': workflow.id}))
  127. except Exception, e:
  128. request.error(_('Could not import workflow: %s' % e))
  129. Workflow.objects.destroy(workflow, request.fs)
  130. raise PopupException(_('Could not import workflow.'), detail=e)
  131. else:
  132. request.error(_('Errors on the form: %s') % workflow_form.errors)
  133. else:
  134. workflow_form = ImportWorkflowForm(instance=workflow)
  135. return render('editor/import_workflow.mako', request, {
  136. 'workflow_form': workflow_form,
  137. 'workflow': workflow,
  138. })
  139. @check_job_access_permission()
  140. def edit_workflow(request, workflow):
  141. history = History.objects.filter(submitter=request.user, job=workflow).order_by('-submission_date')
  142. workflow_form = WorkflowForm(instance=workflow)
  143. user_can_access_job = workflow.is_accessible(request.user)
  144. user_can_edit_job = workflow.is_editable(request.user)
  145. return render('editor/edit_workflow.mako', request, {
  146. 'workflow_form': workflow_form,
  147. 'workflow': workflow,
  148. 'history': history,
  149. 'user_can_access_job': user_can_access_job,
  150. 'user_can_edit_job': user_can_edit_job,
  151. 'job_properties': extract_field_data(workflow_form['job_properties']),
  152. 'link_form': LinkForm(),
  153. 'default_link_form': DefaultLinkForm(action=workflow.start),
  154. 'node_form': NodeForm(),
  155. 'action_forms': [(node_type, design_form_by_type(node_type, request.user, workflow)())
  156. for node_type in ACTION_TYPES.iterkeys()]
  157. })
  158. def delete_workflow(request):
  159. if request.method != 'POST':
  160. raise PopupException(_('A POST request is required.'))
  161. job_ids = request.POST.getlist('job_selection')
  162. for job_id in job_ids:
  163. job = Job.objects.is_accessible_or_exception(request, job_id)
  164. Job.objects.can_edit_or_exception(request, job)
  165. Workflow.objects.destroy(job, request.fs)
  166. request.info(_('Workflow(s) deleted.'))
  167. return redirect(reverse('oozie:list_workflows'))
  168. @check_job_access_permission()
  169. def clone_workflow(request, workflow):
  170. if request.method != 'POST':
  171. raise PopupException(_('A POST request is required.'))
  172. clone = workflow.clone(request.fs, request.user)
  173. response = {'url': reverse('oozie:edit_workflow', kwargs={'workflow': clone.id})}
  174. return HttpResponse(json.dumps(response), mimetype="application/json")
  175. @check_job_access_permission()
  176. def submit_workflow(request, workflow):
  177. ParametersFormSet = formset_factory(ParameterForm, extra=0)
  178. if request.method == 'POST':
  179. params_form = ParametersFormSet(request.POST)
  180. if params_form.is_valid():
  181. mapping = dict([(param['name'], param['value']) for param in params_form.cleaned_data])
  182. job_id = _submit_workflow(request.user, request.fs, workflow, mapping)
  183. request.info(_('Workflow submitted'))
  184. return redirect(reverse('oozie:list_oozie_workflow', kwargs={'job_id': job_id}))
  185. else:
  186. request.error(_('Invalid submission form: %s' % params_form.errors))
  187. else:
  188. parameters = workflow.find_all_parameters()
  189. initial_params = ParameterForm.get_initial_params(dict([(param['name'], param['value']) for param in parameters]))
  190. params_form = ParametersFormSet(initial=initial_params)
  191. popup = render('editor/submit_job_popup.mako', request, {
  192. 'params_form': params_form,
  193. 'action': reverse('oozie:submit_workflow', kwargs={'workflow': workflow.id})
  194. }, force_template=True).content
  195. return HttpResponse(json.dumps(popup), mimetype="application/json")
  196. def _submit_workflow(user, fs, workflow, mapping):
  197. try:
  198. submission = Submission(user, workflow, fs, mapping)
  199. job_id = submission.run()
  200. History.objects.create_from_submission(submission)
  201. return job_id
  202. except RestException, ex:
  203. detail = ex._headers.get('oozie-error-message', ex)
  204. if 'urlopen error' in str(detail):
  205. detail = '%s: %s' % (_('The Oozie server is not running'), detail)
  206. raise PopupException(_("Error submitting workflow %s") % (workflow,), detail=detail)
  207. return redirect(reverse('oozie:list_oozie_workflow', kwargs={'job_id': job_id}))
  208. @check_job_access_permission()
  209. def schedule_workflow(request, workflow):
  210. if Coordinator.objects.filter(workflow=workflow).exists():
  211. request.info(_('You already have some coordinators for this workflow. Submit one or create a new one.'))
  212. return list_coordinators(request, workflow_id=workflow.id)
  213. else:
  214. return create_coordinator(request, workflow=workflow.id)
  215. @check_job_access_permission()
  216. def create_coordinator(request, workflow=None):
  217. if workflow is not None:
  218. coordinator = Coordinator(owner=request.user, schema_version="uri:oozie:coordinator:0.1", workflow=workflow)
  219. else:
  220. coordinator = Coordinator(owner=request.user, schema_version="uri:oozie:coordinator:0.1")
  221. if request.method == 'POST':
  222. coordinator_form = CoordinatorForm(request.POST, instance=coordinator, user=request.user)
  223. if coordinator_form.is_valid():
  224. coordinator = coordinator_form.save()
  225. return redirect(reverse('oozie:edit_coordinator', kwargs={'coordinator': coordinator.id}) + "#step3")
  226. else:
  227. request.error(_('Errors on the form: %s') % coordinator_form.errors)
  228. else:
  229. coordinator_form = CoordinatorForm(instance=coordinator, user=request.user)
  230. return render('editor/create_coordinator.mako', request, {
  231. 'coordinator': coordinator,
  232. 'coordinator_form': coordinator_form,
  233. })
  234. def delete_coordinator(request):
  235. if request.method != 'POST':
  236. raise PopupException(_('A POST request is required.'))
  237. job_ids = request.POST.getlist('job_selection')
  238. for job_id in job_ids:
  239. job = Job.objects.is_accessible_or_exception(request, job_id)
  240. Job.objects.can_edit_or_exception(request, job)
  241. Submission(request.user, job, request.fs, {}).remove_deployment_dir()
  242. job.delete()
  243. request.info(_('Coordinator(s) deleted.'))
  244. return redirect(reverse('oozie:list_coordinators'))
  245. @check_job_access_permission()
  246. @check_job_edition_permission(True)
  247. def edit_coordinator(request, coordinator):
  248. history = History.objects.filter(submitter=request.user, job=coordinator).order_by('-submission_date')
  249. DatasetFormSet = inlineformset_factory(Coordinator, Dataset, form=DatasetForm, max_num=0, can_order=False, can_delete=True)
  250. DataInputFormSet = inlineformset_factory(Coordinator, DataInput, form=DataInputForm, max_num=0, can_order=False, can_delete=True)
  251. DataInputFormSet.form = staticmethod(curry(DataInputForm, coordinator=coordinator))
  252. DataOutputFormSet = inlineformset_factory(Coordinator, DataOutput, form=DataOutputForm, max_num=0, can_order=False, can_delete=True)
  253. DataOutputFormSet.form = staticmethod(curry(DataOutputForm, coordinator=coordinator))
  254. dataset = Dataset(coordinator=coordinator)
  255. dataset_form = DatasetForm(instance=dataset, prefix='create')
  256. NewDataInputFormSet = inlineformset_factory(Coordinator, DataInput, form=DataInputForm, extra=0, can_order=False, can_delete=False)
  257. NewDataInputFormSet.form = staticmethod(curry(DataInputForm, coordinator=coordinator))
  258. NewDataOutputFormSet = inlineformset_factory(Coordinator, DataOutput, form=DataOutputForm, extra=0, can_order=False, can_delete=False)
  259. NewDataOutputFormSet.form = staticmethod(curry(DataOutputForm, coordinator=coordinator))
  260. if request.method == 'POST':
  261. coordinator_form = CoordinatorForm(request.POST, instance=coordinator, user=request.user)
  262. dataset_formset = DatasetFormSet(request.POST, request.FILES, instance=coordinator)
  263. data_input_formset = DataInputFormSet(request.POST, request.FILES, instance=coordinator)
  264. data_output_formset = DataOutputFormSet(request.POST, request.FILES, instance=coordinator)
  265. new_data_input_formset = NewDataInputFormSet(request.POST, request.FILES, instance=coordinator, prefix='input')
  266. new_data_output_formset = NewDataOutputFormSet(request.POST, request.FILES, instance=coordinator, prefix='output')
  267. if coordinator_form.is_valid() and dataset_formset.is_valid() and data_input_formset.is_valid() and data_output_formset.is_valid() \
  268. and new_data_input_formset.is_valid() and new_data_output_formset.is_valid():
  269. coordinator = coordinator_form.save()
  270. dataset_formset.save()
  271. data_input_formset.save()
  272. data_output_formset.save()
  273. new_data_input_formset.save()
  274. new_data_output_formset.save()
  275. request.info(_('Coordinator saved.'))
  276. return redirect(reverse('oozie:edit_coordinator', kwargs={'coordinator': coordinator.id}))
  277. else:
  278. coordinator_form = CoordinatorForm(instance=coordinator, user=request.user)
  279. dataset_formset = DatasetFormSet(instance=coordinator)
  280. data_input_formset = DataInputFormSet(instance=coordinator)
  281. data_output_formset = DataOutputFormSet(instance=coordinator)
  282. new_data_input_formset = NewDataInputFormSet(queryset=DataInput.objects.none(), instance=coordinator, prefix='input')
  283. new_data_output_formset = NewDataOutputFormSet(queryset=DataOutput.objects.none(), instance=coordinator, prefix='output')
  284. return render('editor/edit_coordinator.mako', request, {
  285. 'coordinator': coordinator,
  286. 'coordinator_form': coordinator_form,
  287. 'dataset_formset': dataset_formset,
  288. 'data_input_formset': data_input_formset,
  289. 'data_output_formset': data_output_formset,
  290. 'dataset': dataset,
  291. 'dataset_form': dataset_form,
  292. 'new_data_input_formset': new_data_input_formset,
  293. 'new_data_output_formset': new_data_output_formset,
  294. 'history': history
  295. })
  296. @check_job_access_permission()
  297. @check_job_edition_permission()
  298. def create_coordinator_dataset(request, coordinator):
  299. """Returns {'status' 0/1, data:html or url}"""
  300. dataset = Dataset(coordinator=coordinator)
  301. response = {'status': -1, 'data': 'None'}
  302. if request.method == 'POST':
  303. dataset_form = DatasetForm(request.POST, instance=dataset, prefix='create')
  304. if dataset_form.is_valid():
  305. dataset_form.save()
  306. response['status'] = 0
  307. response['data'] = reverse('oozie:edit_coordinator', kwargs={'coordinator': coordinator.id}) + "#listDataset"
  308. request.info(_('Dataset created'))
  309. else:
  310. ## Bad
  311. response['data'] = _('A POST request is required.')
  312. if response['status'] != 0:
  313. response['data'] = render('editor/create_coordinator_dataset.mako', request, {
  314. 'coordinator': coordinator,
  315. 'dataset_form': dataset_form,
  316. 'dataset': dataset,
  317. }, force_template=True).content
  318. return HttpResponse(json.dumps(response), mimetype="application/json")
  319. @check_dataset_access_permission
  320. @check_dataset_edition_permission()
  321. def edit_coordinator_dataset(request, dataset):
  322. """Returns HTML for modal to edit datasets"""
  323. response = {'status': -1, 'data': 'None'}
  324. if request.method == 'POST':
  325. dataset_form = DatasetForm(request.POST, instance=dataset, prefix='edit')
  326. if dataset_form.is_valid():
  327. dataset = dataset_form.save()
  328. response['status'] = 0
  329. response['data'] = reverse('oozie:edit_coordinator', kwargs={'coordinator': dataset.coordinator.id}) + "#listDataset"
  330. request.info(_('Dataset modified'))
  331. if dataset.start > dataset.coordinator.start:
  332. request.error(_('Beware: dataset start date was after the coordinator start date.'))
  333. else:
  334. response['data'] = dataset_form.errors
  335. else:
  336. dataset_form = DatasetForm(instance=dataset, prefix='edit')
  337. if response['status'] != 0:
  338. response['data'] = render('editor/edit_coordinator_dataset.mako', request, {
  339. 'coordinator': dataset.coordinator,
  340. 'dataset_form': dataset_form,
  341. 'dataset': dataset,
  342. 'path': request.path,
  343. }, force_template=True).content
  344. return HttpResponse(json.dumps(response), mimetype="application/json")
  345. @check_job_access_permission()
  346. @check_job_edition_permission()
  347. def create_coordinator_data(request, coordinator, data_type):
  348. """Returns {'status' 0/1, data:html or url}"""
  349. if data_type == 'input':
  350. data_instance = DataInput(coordinator=coordinator)
  351. DataForm = DataInputForm
  352. else:
  353. data_instance = DataOutput(coordinator=coordinator)
  354. DataForm = DataOutputForm
  355. response = {'status': -1, 'data': 'None'}
  356. if request.method == 'POST':
  357. data_form = DataForm(request.POST, instance=data_instance, coordinator=coordinator, prefix=data_type)
  358. if data_form.is_valid():
  359. data_form.save()
  360. response['status'] = 0
  361. response['data'] = reverse('oozie:edit_coordinator', kwargs={'coordinator': coordinator.id})
  362. request.info(_('Coordinator data created'));
  363. else:
  364. response['data'] = data_form.errors
  365. else:
  366. response['data'] = _('A POST request is required.')
  367. return HttpResponse(json.dumps(response), mimetype="application/json")
  368. @check_job_access_permission()
  369. def clone_coordinator(request, coordinator):
  370. if request.method != 'POST':
  371. raise PopupException(_('A POST request is required.'))
  372. clone = coordinator.clone(request.user)
  373. response = {'url': reverse('oozie:edit_coordinator', kwargs={'coordinator': clone.id})}
  374. return HttpResponse(json.dumps(response), mimetype="application/json")
  375. @check_job_access_permission()
  376. def submit_coordinator(request, coordinator):
  377. ParametersFormSet = formset_factory(ParameterForm, extra=0)
  378. if request.method == 'POST':
  379. params_form = ParametersFormSet(request.POST)
  380. if params_form.is_valid():
  381. mapping = dict([(param['name'], param['value']) for param in params_form.cleaned_data])
  382. job_id = _submit_coordinator(request, coordinator, mapping)
  383. request.info(_('Coordinator submitted.'))
  384. return redirect(reverse('oozie:list_oozie_coordinator', kwargs={'job_id': job_id}))
  385. else:
  386. request.error(_('Invalid submission form: %s' % params_form.errors))
  387. else:
  388. parameters = coordinator.find_all_parameters()
  389. initial_params = ParameterForm.get_initial_params(dict([(param['name'], param['value']) for param in parameters]))
  390. params_form = ParametersFormSet(initial=initial_params)
  391. popup = render('editor/submit_job_popup.mako', request, {
  392. 'params_form': params_form,
  393. 'action': reverse('oozie:submit_coordinator', kwargs={'coordinator': coordinator.id})
  394. }, force_template=True).content
  395. return HttpResponse(json.dumps(popup), mimetype="application/json")
  396. def _submit_coordinator(request, coordinator, mapping):
  397. try:
  398. wf_dir = Submission(request.user, coordinator.workflow, request.fs, mapping).deploy()
  399. properties = {'wf_application_path': request.fs.get_hdfs_path(wf_dir)}
  400. properties.update(mapping)
  401. submission = Submission(request.user, coordinator, request.fs, properties=properties)
  402. job_id = submission.run()
  403. History.objects.create_from_submission(submission)
  404. return job_id
  405. except RestException, ex:
  406. raise PopupException(_("Error submitting coordinator %s") % (coordinator,),
  407. detail=ex._headers.get('oozie-error-message', ex))
  408. def create_bundle(request):
  409. bundle = Bundle(owner=request.user, schema_version='uri:oozie:bundle:0.2')
  410. if request.method == 'POST':
  411. bundle_form = BundleForm(request.POST, instance=bundle)
  412. if bundle_form.is_valid():
  413. bundle = bundle_form.save()
  414. return redirect(reverse('oozie:edit_bundle', kwargs={'bundle': bundle.id}))
  415. else:
  416. request.error(_('Errors on the form: %s') % bundle_form.errors)
  417. else:
  418. bundle_form = BundleForm(instance=bundle)
  419. return render('editor/create_bundle.mako', request, {
  420. 'bundle': bundle,
  421. 'bundle_form': bundle_form,
  422. })
  423. def delete_bundle(request):
  424. if request.method != 'POST':
  425. raise PopupException(_('A POST request is required.'))
  426. job_ids = request.POST.getlist('job_selection')
  427. for job_id in job_ids:
  428. job = Job.objects.is_accessible_or_exception(request, job_id)
  429. Job.objects.can_edit_or_exception(request, job)
  430. Submission(request.user, job, request.fs, {}).remove_deployment_dir()
  431. job.delete()
  432. request.info(_('Bundle(s) deleted.'))
  433. return redirect(reverse('oozie:list_bundles'))
  434. @check_job_access_permission()
  435. @check_job_edition_permission(True)
  436. def edit_bundle(request, bundle):
  437. history = History.objects.filter(submitter=request.user, job=bundle).order_by('-submission_date')
  438. BundledCoordinatorFormSet = inlineformset_factory(Bundle, BundledCoordinator, form=BundledCoordinatorForm, max_num=0, can_order=False, can_delete=True)
  439. bundle_form = BundleForm(instance=bundle)
  440. if request.method == 'POST':
  441. bundle_form = BundleForm(request.POST, instance=bundle)
  442. bundled_coordinator_formset = BundledCoordinatorFormSet(request.POST, instance=bundle)
  443. if bundle_form.is_valid() and bundled_coordinator_formset.is_valid():
  444. bundle = bundle_form.save()
  445. bundled_coordinator_formset.save()
  446. request.info(_('Bundle saved.'))
  447. return redirect(reverse('oozie:list_bundles'))
  448. else:
  449. bundle_form = BundleForm(instance=bundle)
  450. bundled_coordinator_formset = BundledCoordinatorFormSet(instance=bundle)
  451. return render('editor/edit_bundle.mako', request, {
  452. 'bundle': bundle,
  453. 'bundle_form': bundle_form,
  454. 'bundled_coordinator_formset': bundled_coordinator_formset,
  455. 'bundled_coordinator_html_form': get_create_bundled_coordinator_html(request, bundle),
  456. 'history': history
  457. })
  458. @check_job_access_permission()
  459. @check_job_edition_permission(True)
  460. def create_bundled_coordinator(request, bundle):
  461. bundled_coordinator_instance = BundledCoordinator(bundle=bundle)
  462. response = {'status': -1, 'data': 'None'}
  463. if request.method == 'POST':
  464. bundled_coordinator_form = BundledCoordinatorForm(request.POST, instance=bundled_coordinator_instance, prefix='create-bundled-coordinator')
  465. if bundled_coordinator_form.is_valid():
  466. bundled_coordinator_form.save()
  467. response['status'] = 0
  468. response['data'] = reverse('oozie:edit_bundle', kwargs={'bundle': bundle.id}) + "#listCoordinators"
  469. request.info(_('Coordinator added to the bundle!'))
  470. else:
  471. bundled_coordinator_form = BundledCoordinatorForm(instance=bundled_coordinator_instance, prefix='create-bundled-coordinator')
  472. if response['status'] != 0:
  473. response['data'] = get_create_bundled_coordinator_html(request, bundle, bundled_coordinator_form=bundled_coordinator_form)
  474. return HttpResponse(json.dumps(response), mimetype="application/json")
  475. def get_create_bundled_coordinator_html(request, bundle, bundled_coordinator_form=None):
  476. if bundled_coordinator_form is None:
  477. bundled_coordinator_instance = BundledCoordinator(bundle=bundle)
  478. bundled_coordinator_form = BundledCoordinatorForm(instance=bundled_coordinator_instance, prefix='create-bundled-coordinator')
  479. return render('editor/create_bundled_coordinator.mako', request, {
  480. 'bundle': bundle,
  481. 'bundled_coordinator_form': bundled_coordinator_form,
  482. }, force_template=True).content
  483. @check_job_access_permission()
  484. @check_job_edition_permission(True)
  485. def edit_bundled_coordinator(request, bundle, bundled_coordinator):
  486. bundled_coordinator_instance = BundledCoordinator.objects.get(id=bundled_coordinator) # todo secu
  487. response = {'status': -1, 'data': 'None'}
  488. if request.method == 'POST':
  489. bundled_coordinator_form = BundledCoordinatorForm(request.POST, instance=bundled_coordinator_instance, prefix='edit-bundled-coordinator')
  490. if bundled_coordinator_form.is_valid():
  491. bundled_coordinator_form.save()
  492. response['status'] = 0
  493. response['data'] = reverse('oozie:edit_bundle', kwargs={'bundle': bundle.id}) + "#listCoordinators"
  494. request.info(_('Bundled coordinator updated!'))
  495. else:
  496. bundled_coordinator_form = BundledCoordinatorForm(instance=bundled_coordinator_instance, prefix='edit-bundled-coordinator')
  497. if response['status'] != 0:
  498. response['data'] = render('editor/edit_bundled_coordinator.mako', request, {
  499. 'bundle': bundle,
  500. 'bundled_coordinator_form': bundled_coordinator_form,
  501. 'bundled_coordinator_instance': bundled_coordinator_instance,
  502. }, force_template=True).content
  503. return HttpResponse(json.dumps(response), mimetype="application/json")
  504. @check_job_access_permission()
  505. def clone_bundle(request, bundle):
  506. if request.method != 'POST':
  507. raise PopupException(_('A POST request is required.'))
  508. clone = bundle.clone(request.user)
  509. response = {'url': reverse('oozie:edit_bundle', kwargs={'bundle': clone.id})}
  510. return HttpResponse(json.dumps(response), mimetype="application/json")
  511. @check_job_access_permission()
  512. def submit_bundle(request, bundle):
  513. ParametersFormSet = formset_factory(ParameterForm, extra=0)
  514. if request.method == 'POST':
  515. params_form = ParametersFormSet(request.POST)
  516. if params_form.is_valid():
  517. mapping = dict([(param['name'], param['value']) for param in params_form.cleaned_data])
  518. job_id = _submit_bundle(request, bundle, mapping)
  519. request.info(_('Bundle submitted.'))
  520. return redirect(reverse('oozie:list_oozie_bundle', kwargs={'job_id': job_id}))
  521. else:
  522. request.error(_('Invalid submission form: %s' % params_form.errors))
  523. else:
  524. parameters = bundle.find_all_parameters()
  525. initial_params = ParameterForm.get_initial_params(dict([(param['name'], param['value']) for param in parameters]))
  526. params_form = ParametersFormSet(initial=initial_params)
  527. popup = render('editor/submit_job_popup.mako', request, {
  528. 'params_form': params_form,
  529. 'action': reverse('oozie:submit_bundle', kwargs={'bundle': bundle.id})
  530. }, force_template=True).content
  531. return HttpResponse(json.dumps(popup), mimetype="application/json")
  532. def _submit_bundle(request, bundle, properties):
  533. try:
  534. deployment_dirs = {}
  535. for bundled in bundle.coordinators.all():
  536. wf_dir = Submission(request.user, bundled.coordinator.workflow, request.fs, properties).deploy()
  537. deployment_dirs['wf_%s_dir' % bundled.coordinator.workflow.id] = request.fs.get_hdfs_path(wf_dir)
  538. coord_dir = Submission(request.user, bundled.coordinator, request.fs, properties).deploy()
  539. deployment_dirs['coord_%s_dir' % bundled.coordinator.id] = coord_dir
  540. properties.update(deployment_dirs)
  541. submission = Submission(request.user, bundle, request.fs, properties=properties)
  542. job_id = submission.run()
  543. History.objects.create_from_submission(submission)
  544. return job_id
  545. except RestException, ex:
  546. raise PopupException(_("Error submitting bundle %s") % (bundle,),
  547. detail=ex._headers.get('oozie-error-message', ex))
  548. def list_history(request):
  549. """
  550. List the job submission history.
  551. Normal users can only look at their own submissions.
  552. """
  553. history = History.objects
  554. if not request.user.is_superuser:
  555. history = history.filter(submitter=request.user)
  556. history = history.order_by('-submission_date')
  557. return render('editor/list_history.mako', request, {
  558. 'history': history,
  559. })
  560. def list_history_record(request, record_id):
  561. """
  562. List a job submission history.
  563. Normal users can only look at their own jobs.
  564. """
  565. history = History.objects
  566. if not request.user.is_superuser:
  567. history.filter(submitter=request.user)
  568. history = history.get(id=record_id)
  569. return render('editor/list_history_record.mako', request, {
  570. 'record': history,
  571. })
  572. def setup_app(request):
  573. if request.method != 'POST':
  574. raise PopupException(_('A POST request is required.'))
  575. try:
  576. oozie_setup.Command().handle_noargs()
  577. activate_translation(request.LANGUAGE_CODE)
  578. request.info(_('Workspaces and examples installed.'))
  579. except WebHdfsException, e:
  580. raise PopupException(_('The app setup could complete.'), detail=e)
  581. return redirect(reverse('oozie:list_workflows'))
  582. def jasmine(request):
  583. return render('editor/jasmine.mako', request, None)