Skip to content

Commit 91d52f1

Browse files
committed
more o11y
1 parent 22577c2 commit 91d52f1

4 files changed

Lines changed: 32 additions & 2 deletions

File tree

apps/webapp/app/services/runsReplicationService.server.ts

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1307,6 +1307,7 @@ export class RunsReplicationService {
13071307
run.id, // run_id
13081308
run.updatedAt.getTime(), // updated_at
13091309
run.createdAt.getTime(), // created_at
1310+
run.queueTimestamp?.getTime() ?? null, // queue_timestamp
13101311
run.status, // status
13111312
environmentType, // environment_type
13121313
run.friendlyId, // friendly_id

apps/webapp/test/runsReplicationService.part1.test.ts

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -73,6 +73,7 @@ describe("RunsReplicationService (part 1/7)", () => {
7373
},
7474
});
7575

76+
const queueTimestamp = new Date("2026-08-11T12:34:56.789Z");
7677
const taskRun = await prisma.taskRun.create({
7778
data: {
7879
friendlyId: "run_1234",
@@ -81,6 +82,7 @@ describe("RunsReplicationService (part 1/7)", () => {
8182
traceId: "1234",
8283
spanId: "1234",
8384
queue: "test",
85+
queueTimestamp,
8486
workerQueue: "us-east-1-next",
8587
region: "us-east-1",
8688
planType: "free",
@@ -100,7 +102,8 @@ describe("RunsReplicationService (part 1/7)", () => {
100102

101103
const queryRuns = clickhouse.reader.query({
102104
name: "runs-replication",
103-
query: "SELECT * FROM trigger_dev.task_runs_v2",
105+
query:
106+
"SELECT *, toString(toUnixTimestamp64Milli(queue_timestamp)) AS queue_timestamp_ms FROM trigger_dev.task_runs_v2",
104107
schema: z.any(),
105108
});
106109

@@ -125,6 +128,7 @@ describe("RunsReplicationService (part 1/7)", () => {
125128
organization_id: organization.id,
126129
environment_type: "DEVELOPMENT",
127130
engine: "V2",
131+
queue_timestamp_ms: queueTimestamp.getTime().toString(),
128132
trigger_source: "api",
129133
root_trigger_source: "dashboard",
130134
is_warm_start: 1,

internal-packages/run-engine/src/engine/systems/dequeueSystem.ts

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -156,6 +156,10 @@ export class DequeueSystem {
156156

157157
const orgId = message.message.orgId;
158158
const runId = message.messageId;
159+
const queueWaitMs =
160+
typeof message.message.eligibleAtMs === "number"
161+
? Math.max(0, Date.now() - message.message.eligibleAtMs)
162+
: undefined;
159163

160164
this.$.logger.info("DequeueSystem.dequeueFromWorkerQueue dequeued message", {
161165
runId,
@@ -174,6 +178,9 @@ export class DequeueSystem {
174178
span.setAttribute("consumer_id", consumerId);
175179
span.setAttribute("worker_queue", workerQueue);
176180
span.setAttribute("blocking_pop", blockingPop ?? true);
181+
if (queueWaitMs !== undefined) {
182+
span.setAttribute("queue_wait_ms", queueWaitMs);
183+
}
177184

178185
//lock the run so nothing else can modify it
179186
try {

internal-packages/schedule-engine/src/engine/index.ts

Lines changed: 19 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -225,18 +225,31 @@ export class ScheduleEngine {
225225
instance.taskSchedule.timezone,
226226
nominalAt
227227
);
228-
const { effectiveAt: candidateEffectiveAt } = calculateEffectiveScheduleTime({
228+
const {
229+
effectiveAt: candidateEffectiveAt,
230+
effectiveRangeMs,
231+
windowMs,
232+
offsetMs: candidateDelayMs,
233+
rangeWasClamped,
234+
} = calculateEffectiveScheduleTime({
229235
nominalAt,
230236
nextNominalAt,
231237
schedulePhase,
232238
window: scheduleWindow,
233239
});
234240
const effectiveAt = this.options.cronSpreadEnabled ? candidateEffectiveAt : nominalAt;
241+
const appliedDelayMs = effectiveAt.getTime() - nominalAt.getTime();
235242

236243
span.setAttribute("cron_spread_enabled", this.options.cronSpreadEnabled);
244+
span.setAttribute("schedule_window_type", scheduleWindow?.type ?? "none");
237245
span.setAttribute("next_scheduled_timestamp", nominalAt.toISOString());
238246
span.setAttribute("candidate_effective_schedule_time", candidateEffectiveAt.toISOString());
239247
span.setAttribute("effective_schedule_time", effectiveAt.toISOString());
248+
span.setAttribute("candidate_delay_ms", candidateDelayMs);
249+
span.setAttribute("applied_delay_ms", appliedDelayMs);
250+
span.setAttribute("schedule_window_ms", windowMs);
251+
span.setAttribute("effective_range_ms", effectiveRangeMs);
252+
span.setAttribute("schedule_range_was_clamped", rangeWasClamped);
240253

241254
const schedulingDelayMs = effectiveAt.getTime() - Date.now();
242255
span.setAttribute("scheduling_delay_ms", schedulingDelayMs);
@@ -248,6 +261,11 @@ export class ScheduleEngine {
248261
candidateEffectiveAt: candidateEffectiveAt.toISOString(),
249262
effectiveAt: effectiveAt.toISOString(),
250263
cronSpreadEnabled: this.options.cronSpreadEnabled,
264+
scheduleWindowType: scheduleWindow?.type ?? "none",
265+
candidateDelayMs,
266+
appliedDelayMs,
267+
effectiveRangeMs,
268+
rangeWasClamped,
251269
schedulingDelayMs,
252270
generatorExpression: instance.taskSchedule.generatorExpression,
253271
timezone: instance.taskSchedule.timezone,

0 commit comments

Comments
 (0)