Skip to content
GitLab
Explore
Sign in
Primary navigation
Search or go to…
Project
V
vospace-rest
Manage
Activity
Members
Labels
Plan
Issues
Issue boards
Milestones
Wiki
Redmine
Code
Merge requests
Repository
Branches
Commits
Tags
Repository graph
Compare revisions
Snippets
Build
Pipelines
Jobs
Pipeline schedules
Artifacts
Deploy
Releases
Package registry
Container registry
Model registry
Operate
Terraform modules
Analyze
Value stream analytics
Contributor analytics
CI/CD analytics
Repository analytics
Model experiments
Help
Help
Support
GitLab documentation
Compare GitLab plans
Community forum
Contribute to GitLab
Provide feedback
Keyboard shortcuts
?
Snippets
Groups
Projects
Show more breadcrumbs
VOSpace INAF
vospace-rest
Commits
b58a6eac
Commit
b58a6eac
authored
4 years ago
by
Sonia Zorba
Browse files
Options
Downloads
Patches
Plain Diff
Set JobPhase 'QUEUED' to async recall jobs startup
parent
042ca6ad
No related branches found
No related tags found
No related merge requests found
Changes
2
Show whitespace changes
Inline
Side-by-side
Showing
2 changed files
src/main/java/it/inaf/oats/vospace/JobService.java
+18
-4
18 additions, 4 deletions
src/main/java/it/inaf/oats/vospace/JobService.java
src/test/java/it/inaf/oats/vospace/JobServiceTest.java
+113
-0
113 additions, 0 deletions
src/test/java/it/inaf/oats/vospace/JobServiceTest.java
with
131 additions
and
4 deletions
src/main/java/it/inaf/oats/vospace/JobService.java
+
18
−
4
View file @
b58a6eac
...
@@ -10,10 +10,14 @@ import it.inaf.oats.vospace.exception.ErrorSummaryFactory;
...
@@ -10,10 +10,14 @@ import it.inaf.oats.vospace.exception.ErrorSummaryFactory;
import
org.springframework.beans.factory.annotation.Autowired
;
import
org.springframework.beans.factory.annotation.Autowired
;
import
org.springframework.stereotype.Service
;
import
org.springframework.stereotype.Service
;
import
it.inaf.oats.vospace.exception.VoSpaceErrorSummarizableException
;
import
it.inaf.oats.vospace.exception.VoSpaceErrorSummarizableException
;
import
org.slf4j.Logger
;
import
org.slf4j.LoggerFactory
;
@Service
@Service
public
class
JobService
{
public
class
JobService
{
private
static
final
Logger
LOG
=
LoggerFactory
.
getLogger
(
JobService
.
class
);
@Autowired
@Autowired
private
JobDAO
jobDAO
;
private
JobDAO
jobDAO
;
...
@@ -33,6 +37,8 @@ public class JobService {
...
@@ -33,6 +37,8 @@ public class JobService {
public
void
setJobPhase
(
JobSummary
job
,
String
phase
)
{
public
void
setJobPhase
(
JobSummary
job
,
String
phase
)
{
LOG
.
trace
(
"Job "
+
job
.
getJobId
()
+
" phase set to "
+
phase
);
// TODO: check allowed job phase transitions
// TODO: check allowed job phase transitions
switch
(
phase
)
{
switch
(
phase
)
{
case
"RUN"
:
case
"RUN"
:
...
@@ -48,11 +54,19 @@ public class JobService {
...
@@ -48,11 +54,19 @@ public class JobService {
private
void
startJob
(
JobSummary
job
)
{
private
void
startJob
(
JobSummary
job
)
{
try
{
try
{
job
.
setPhase
(
ExecutionPhase
.
EXECUTING
);
jobDAO
.
updateJob
(
job
);
Transfer
transfer
=
uriService
.
getTransfer
(
job
);
Transfer
transfer
=
uriService
.
getTransfer
(
job
);
ExecutionPhase
phase
;
if
(
transfer
.
getProtocols
().
stream
().
anyMatch
(
p
->
"ia2:async-recall"
.
equals
(
p
.
getUri
())))
{
// Async recall from tape jobs are queued. They will be started by VOSpace transfer service
phase
=
ExecutionPhase
.
QUEUED
;
}
else
{
phase
=
ExecutionPhase
.
EXECUTING
;
}
job
.
setPhase
(
phase
);
jobDAO
.
updateJob
(
job
);
switch
(
getJobType
(
transfer
))
{
switch
(
getJobType
(
transfer
))
{
case
pullToVoSpace:
case
pullToVoSpace:
handlePullToVoSpace
(
job
,
transfer
);
handlePullToVoSpace
(
job
,
transfer
);
...
...
This diff is collapsed.
Click to expand it.
src/test/java/it/inaf/oats/vospace/JobServiceTest.java
0 → 100644
+
113
−
0
View file @
b58a6eac
package
it.inaf.oats.vospace
;
import
it.inaf.oats.vospace.exception.NodeBusyException
;
import
it.inaf.oats.vospace.persistence.JobDAO
;
import
net.ivoa.xml.uws.v1.ExecutionPhase
;
import
net.ivoa.xml.uws.v1.JobSummary
;
import
net.ivoa.xml.vospace.v2.Protocol
;
import
net.ivoa.xml.vospace.v2.Transfer
;
import
org.junit.jupiter.api.Test
;
import
org.junit.jupiter.api.extension.ExtendWith
;
import
static
org
.
mockito
.
ArgumentMatchers
.
any
;
import
static
org
.
mockito
.
ArgumentMatchers
.
argThat
;
import
org.mockito.InjectMocks
;
import
org.mockito.Mock
;
import
static
org
.
mockito
.
Mockito
.
doThrow
;
import
static
org
.
mockito
.
Mockito
.
times
;
import
static
org
.
mockito
.
Mockito
.
verify
;
import
static
org
.
mockito
.
Mockito
.
when
;
import
org.mockito.junit.jupiter.MockitoExtension
;
@ExtendWith
(
MockitoExtension
.
class
)
public
class
JobServiceTest
{
@Mock
private
JobDAO
jobDAO
;
@Mock
private
UriService
uriService
;
@Mock
private
AsyncTransferService
asyncTransfService
;
@InjectMocks
private
JobService
jobService
;
@Test
public
void
testStartJobDefault
()
{
when
(
uriService
.
getTransfer
(
any
())).
thenReturn
(
getHttpTransfer
());
JobSummary
job
=
new
JobSummary
();
jobService
.
setJobPhase
(
job
,
"RUN"
);
verify
(
jobDAO
,
times
(
1
)).
updateJob
(
argThat
(
j
->
ExecutionPhase
.
EXECUTING
.
equals
(
j
.
getPhase
())));
}
@Test
public
void
testStartJobTape
()
{
when
(
uriService
.
getTransfer
(
any
())).
thenReturn
(
getTapeTransfer
());
JobSummary
job
=
new
JobSummary
();
jobService
.
setJobPhase
(
job
,
"RUN"
);
verify
(
jobDAO
,
times
(
1
)).
updateJob
(
argThat
(
j
->
ExecutionPhase
.
QUEUED
.
equals
(
j
.
getPhase
())));
}
@Test
public
void
testStartJobVoSpaceError
()
{
when
(
uriService
.
getTransfer
(
any
())).
thenReturn
(
getTapeTransfer
());
when
(
asyncTransfService
.
startJob
(
any
())).
thenThrow
(
new
NodeBusyException
(
"/foo"
));
JobSummary
job
=
new
JobSummary
();
jobService
.
setJobPhase
(
job
,
"RUN"
);
verify
(
jobDAO
,
times
(
2
)).
updateJob
(
argThat
(
j
->
ExecutionPhase
.
ERROR
.
equals
(
j
.
getPhase
())));
}
@Test
public
void
testStartJobUnexpectedError
()
{
when
(
uriService
.
getTransfer
(
any
())).
thenReturn
(
getTapeTransfer
());
when
(
asyncTransfService
.
startJob
(
any
())).
thenThrow
(
new
NullPointerException
());
JobSummary
job
=
new
JobSummary
();
jobService
.
setJobPhase
(
job
,
"RUN"
);
verify
(
jobDAO
,
times
(
2
)).
updateJob
(
argThat
(
j
->
ExecutionPhase
.
ERROR
.
equals
(
j
.
getPhase
())));
}
@Test
public
void
testSyncJobResultVoSpaceError
()
{
when
(
uriService
.
getTransfer
(
any
())).
thenReturn
(
getHttpTransfer
());
doThrow
(
new
NodeBusyException
(
"/foo"
)).
when
(
uriService
).
setSyncTransferEndpoints
(
any
());
jobService
.
createSyncJobResult
(
new
JobSummary
());
verify
(
jobDAO
,
times
(
1
)).
createJob
(
argThat
(
j
->
ExecutionPhase
.
ERROR
.
equals
(
j
.
getPhase
())));
}
@Test
public
void
testSyncJobResultUnexpectedError
()
{
when
(
uriService
.
getTransfer
(
any
())).
thenReturn
(
getHttpTransfer
());
doThrow
(
new
NullPointerException
()).
when
(
uriService
).
setSyncTransferEndpoints
(
any
());
jobService
.
createSyncJobResult
(
new
JobSummary
());
verify
(
jobDAO
,
times
(
1
)).
createJob
(
argThat
(
j
->
ExecutionPhase
.
ERROR
.
equals
(
j
.
getPhase
())));
}
private
Transfer
getHttpTransfer
()
{
Transfer
transfer
=
new
Transfer
();
transfer
.
setDirection
(
"pullFromVoSpace"
);
Protocol
protocol
=
new
Protocol
();
protocol
.
setUri
(
"ivo://ivoa.net/vospace/core#httpget"
);
transfer
.
getProtocols
().
add
(
protocol
);
return
transfer
;
}
private
Transfer
getTapeTransfer
()
{
Transfer
transfer
=
new
Transfer
();
transfer
.
setDirection
(
"pullToVoSpace"
);
Protocol
protocol
=
new
Protocol
();
protocol
.
setUri
(
"ia2:async-recall"
);
transfer
.
getProtocols
().
add
(
protocol
);
return
transfer
;
}
}
This diff is collapsed.
Click to expand it.
Preview
0%
Loading
Try again
or
attach a new file
.
Cancel
You are about to add
0
people
to the discussion. Proceed with caution.
Finish editing this message first!
Save comment
Cancel
Please
register
or
sign in
to comment