pr updates
This commit is contained in:
parent
5737b73949
commit
2cb071d9fb
|
@ -27,9 +27,6 @@ export class StartJobs {
|
|||
|
||||
for (const job of jobs) {
|
||||
if (!this.hasCapacity()) {
|
||||
this.logger.info(
|
||||
`[run] Max concurrent jobs reached (${this.options.maxConcurrentJobs}), stopping`
|
||||
);
|
||||
break;
|
||||
}
|
||||
|
||||
|
@ -56,7 +53,7 @@ export class StartJobs {
|
|||
}
|
||||
|
||||
private getCurrentExecutions(): JobExecution[] {
|
||||
return Array.from(this.runnables.values()).map((r) => r.exec);
|
||||
return Array.from(this.runnables.values()).map((runner) => runner.exec);
|
||||
}
|
||||
|
||||
private hasCapacity(): boolean {
|
||||
|
@ -69,6 +66,18 @@ export class StartJobs {
|
|||
return available;
|
||||
}
|
||||
|
||||
private async tryJobExecution(job: JobDefinition): Promise<JobExecution> {
|
||||
const handlers = await this.jobRepository.getHandlers(job);
|
||||
if (handlers.length === 0) {
|
||||
this.logger.error(`[run] No handlers for job ${job.id}`);
|
||||
throw new Error("No handlers for job");
|
||||
}
|
||||
|
||||
const runnable = this.jobRepository.getRunnableJob(job);
|
||||
|
||||
return this.trackExecution(job, handlers, runnable);
|
||||
}
|
||||
|
||||
private async trackExecution(
|
||||
job: JobDefinition,
|
||||
handlers: Handler[],
|
||||
|
@ -93,16 +102,4 @@ export class StartJobs {
|
|||
|
||||
return jobExec;
|
||||
}
|
||||
|
||||
private async tryJobExecution(job: JobDefinition): Promise<JobExecution> {
|
||||
const handlers = await this.jobRepository.getHandlers(job);
|
||||
if (handlers.length === 0) {
|
||||
this.logger.error(`[run] No handlers for job ${job.id}`);
|
||||
throw new Error("No handlers for job");
|
||||
}
|
||||
|
||||
const runnable = this.jobRepository.getRunnableJob(job);
|
||||
|
||||
return this.trackExecution(job, handlers, runnable);
|
||||
}
|
||||
}
|
||||
|
|
|
@ -16,7 +16,7 @@ import {
|
|||
PostgresJobExecutionRepository,
|
||||
InMemoryJobExecutionRepository,
|
||||
ArbitrumEvmJsonRPCBlockRepository,
|
||||
} from "./";
|
||||
} from ".";
|
||||
import { HttpClient } from "../rpc/http/HttpClient";
|
||||
import {
|
||||
JobExecutionRepository,
|
||||
|
|
|
@ -16,7 +16,6 @@ import {
|
|||
SolanaSlotRepository,
|
||||
StatRepository,
|
||||
} from "../../../domain/repositories";
|
||||
import log from "../../log";
|
||||
import { FileMetadataRepository, SnsEventRepository } from "..";
|
||||
import {
|
||||
solanaLogMessagePublishedMapper,
|
||||
|
@ -25,6 +24,7 @@ import {
|
|||
evmStandardRelayDelivered,
|
||||
evmTransferRedeemedMapper,
|
||||
} from "../../mappers";
|
||||
import log from "../../log";
|
||||
|
||||
export class StaticJobRepository implements JobRepository {
|
||||
private fileRepo: FileMetadataRepository;
|
||||
|
|
Loading…
Reference in New Issue