File
Blob: src/worker/tasks/context.ts
| 1 | import type { CacheContext } from "@/worker/cache"; |
| 2 | import type { Logger, LoggerContext } from "@/worker/common/logger"; |
| 3 | import type { Db } from "@/worker/db/d1/client"; |
| 4 | |
| 5 | import { createLogger } from "@/worker/common"; |
| 6 | import { createDb } from "@/worker/db/d1/client"; |
| 7 | import { |
| 8 | MAX_SIMULTANEOUS_CONNECTIONS, |
| 9 | SubrequestLimiter, |
| 10 | countSubrequest, |
| 11 | } from "@/worker/git/operations/limits"; |
| 12 | |
| 13 | export type QueueLogContext = Omit<LoggerContext, "requestId">; |
| 14 | |
| 15 | export type QueueTaskContext = { |
| 16 | db: Db; |
| 17 | cacheCtx: CacheContext; |
| 18 | limiter: SubrequestLimiter; |
| 19 | logFor: (context: QueueLogContext) => Logger; |
| 20 | }; |
| 21 | |
| 22 | export function createQueueTaskContext(args: { |
| 23 | env: Env; |
| 24 | ctx: ExecutionContext; |
| 25 | repoLabel: string; |
| 26 | operation: string; |
| 27 | subrequestBudget: number; |
| 28 | }): QueueTaskContext { |
| 29 | const limiter = new SubrequestLimiter(MAX_SIMULTANEOUS_CONNECTIONS); |
| 30 | return { |
| 31 | db: createDb(args.env.DB), |
| 32 | cacheCtx: { |
| 33 | req: new Request( |
| 34 | `https://queue.internal/${encodeURIComponent(args.repoLabel)}/${args.operation}` |
| 35 | ), |
| 36 | ctx: args.ctx, |
| 37 | memo: { |
| 38 | repoId: args.repoLabel, |
| 39 | limiter, |
| 40 | subreqBudget: args.subrequestBudget, |
| 41 | }, |
| 42 | }, |
| 43 | limiter, |
| 44 | logFor: (context) => createLogger(args.env.LOG_LEVEL, context), |
| 45 | }; |
| 46 | } |
| 47 | |
| 48 | export function retryQueueMessage( |
| 49 | message: { retry: (options?: { delaySeconds?: number }) => void }, |
| 50 | seconds: number |
| 51 | ): void { |
| 52 | message.retry({ delaySeconds: seconds }); |
| 53 | } |
| 54 | |
| 55 | export function logSoftBudgetExhausted(args: { |
| 56 | cacheCtx: CacheContext; |
| 57 | log: Logger; |
| 58 | flagPrefix: string; |
| 59 | op: string; |
| 60 | count?: number; |
| 61 | }): void { |
| 62 | if (countSubrequest(args.cacheCtx, args.count)) return; |
| 63 | args.cacheCtx.memo = args.cacheCtx.memo || {}; |
| 64 | args.cacheCtx.memo.flags = args.cacheCtx.memo.flags || new Set(); |
| 65 | const flag = `${args.flagPrefix}:${args.op}`; |
| 66 | if (args.cacheCtx.memo.flags.has(flag)) return; |
| 67 | args.cacheCtx.memo.flags.add(flag); |
| 68 | args.log.warn("soft-budget-exhausted", { op: args.op }); |
| 69 | } |