-
-
Notifications
You must be signed in to change notification settings - Fork 80.8k
Expand file tree
/
Copy pathwebhook-ack.ts
More file actions
143 lines (133 loc) · 4.64 KB
/
Copy pathwebhook-ack.ts
File metadata and controls
143 lines (133 loc) · 4.64 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
// Line plugin module implements webhook acknowledgement helpers.
import type { webhook } from "@line/bot-sdk";
import { danger, type RuntimeEnv } from "openclaw/plugin-sdk/runtime-env";
// LINE classifies responses after 2s as request_timeout. Fail before that
// limit to release ingress capacity; this deadline never acknowledges an event.
export const LINE_WEBHOOK_RESPONSE_DEADLINE_MS = 1_500;
function createLineWebhookResponseDeadlineError(): Error {
return new Error("LINE webhook response deadline elapsed before acceptance");
}
export function assertLineWebhookResponseDeadline(responseDeadlineAt: number): void {
if (Date.now() >= responseDeadlineAt) {
throw createLineWebhookResponseDeadlineError();
}
}
export type LineWebhookDispatchCallbacks = {
/** Called once when one webhook event is durably owned or needs no dispatch. */
onEventAccepted: (event: webhook.Event) => void | Promise<void>;
/** Aborted when ingress can no longer acknowledge pre-acceptance work safely. */
abortSignal: AbortSignal;
};
export type LineWebhookDispatchHandler = (body: webhook.CallbackRequest) => Promise<void>;
export type LineWebhookAcceptanceDispatchHandler = (
body: webhook.CallbackRequest,
callbacks: LineWebhookDispatchCallbacks,
) => Promise<void>;
function logLineWebhookDispatchError(runtime: RuntimeEnv | undefined, err: unknown): void {
runtime?.error?.(danger(`line webhook dispatch failed: ${String(err)}`));
}
export function dispatchLineWebhookInBackground(params: {
body: webhook.CallbackRequest;
dispatch: LineWebhookDispatchHandler;
runtime?: RuntimeEnv;
}): void {
void Promise.resolve()
.then(() => params.dispatch(params.body))
.catch((error: unknown) => logLineWebhookDispatchError(params.runtime, error));
}
/**
* Wait until every callback event is either handed to the durable reply lane
* or completes without one. A completed legacy handler is safe fallback only
* because all of its events have already settled.
*/
export async function waitForLineWebhookDispatchAcceptance(params: {
body: webhook.CallbackRequest;
dispatch: LineWebhookAcceptanceDispatchHandler;
responseDeadlineAt?: number;
runtime?: RuntimeEnv;
}): Promise<void> {
const expectedEvents = new Set(params.body.events ?? []);
if (expectedEvents.size === 0) {
return;
}
const responseDeadlineAt =
params.responseDeadlineAt ?? Date.now() + LINE_WEBHOOK_RESPONSE_DEADLINE_MS;
assertLineWebhookResponseDeadline(responseDeadlineAt);
const acceptedEvents = new Set<webhook.Event>();
const dispatchAbort = new AbortController();
let settled = false;
let acceptanceDeadline: ReturnType<typeof setTimeout> | undefined;
let resolveAcceptance!: () => void;
let rejectAcceptancePromise!: (error: unknown) => void;
const acceptance = new Promise<void>((resolve, reject) => {
resolveAcceptance = resolve;
rejectAcceptancePromise = reject;
});
const clearAcceptanceDeadline = () => {
if (acceptanceDeadline) {
clearTimeout(acceptanceDeadline);
acceptanceDeadline = undefined;
}
};
const rejectAcceptance = (error: unknown) => {
if (settled) {
return;
}
settled = true;
dispatchAbort.abort(error);
clearAcceptanceDeadline();
rejectAcceptancePromise(error);
};
const acceptRemainingEvents = () => {
if (settled) {
return;
}
try {
assertLineWebhookResponseDeadline(responseDeadlineAt);
} catch (error) {
rejectAcceptance(error);
return;
}
settled = true;
clearAcceptanceDeadline();
resolveAcceptance();
};
acceptanceDeadline = setTimeout(
() => rejectAcceptance(createLineWebhookResponseDeadlineError()),
Math.max(0, responseDeadlineAt - Date.now()),
);
const dispatch = Promise.resolve().then(() => {
assertLineWebhookResponseDeadline(responseDeadlineAt);
return params.dispatch(params.body, {
abortSignal: dispatchAbort.signal,
onEventAccepted: (event) => {
if (settled || !expectedEvents.has(event)) {
return;
}
try {
assertLineWebhookResponseDeadline(responseDeadlineAt);
} catch (error) {
rejectAcceptance(error);
return;
}
acceptedEvents.add(event);
if (acceptedEvents.size === expectedEvents.size) {
settled = true;
clearAcceptanceDeadline();
resolveAcceptance();
}
},
});
});
void dispatch.then(
() => acceptRemainingEvents(),
(error: unknown) => {
const acceptedBeforeFailure = settled;
rejectAcceptance(error);
if (acceptedBeforeFailure) {
logLineWebhookDispatchError(params.runtime, error);
}
},
);
await acceptance;
}