惯性聚合 高效追踪和阅读你感兴趣的博客、新闻、科技资讯
阅读原文 在惯性聚合中打开

推荐订阅源

让小产品的独立变现更简单 - ezindie.com
让小产品的独立变现更简单 - ezindie.com
U
Unit 42
IT之家
IT之家
Y
Y Combinator Blog
T
Tailwind CSS Blog
B
Blog
大猫的无限游戏
大猫的无限游戏
博客园 - 叶小钗
Jina AI
Jina AI
freeCodeCamp Programming Tutorials: Python, JavaScript, Git & More
I
InfoQ
J
Java Code Geeks
F
Fortinet All Blogs
T
The Blog of Author Tim Ferriss
OSCHINA 社区最新新闻
OSCHINA 社区最新新闻
H
Hackread – Cybersecurity News, Data Breaches, AI and More
人人都是产品经理
人人都是产品经理
腾讯CDC
Hugging Face - Blog
Hugging Face - Blog
GbyAI
GbyAI
博客园 - 司徒正美
The GitHub Blog
The GitHub Blog
酷 壳 – CoolShell
酷 壳 – CoolShell
L
LangChain Blog

Recent Commits to openclaw:main

test: merge chat side-result checks · openclaw/openclaw@ddd2c2a test: merge cron history checks · openclaw/openclaw@f7eb746 test: merge responsive navigation shell checks · openclaw/openclaw@c2e4b47 docs(changelog): add codex oauth fixes · openclaw/openclaw@628e6cd test: merge navigation routing cases · openclaw/openclaw@5d8cecb Tests: mock channel registry bundled fallback · openclaw/openclaw@2b08233 Secrets: avoid broad web search discovery for single plugin config · openclaw/openclaw@a464f59 test: merge config view browser checks · openclaw/openclaw@20cf511 fix(status): align oauth health with runtime · openclaw/openclaw@eed7116 feat: add macOS screen snapshots for monitor preview (#67954) thanks … · openclaw/openclaw@f377db1 fix: report shared auth scopes in hello-ok (#67810) thanks @BunsDev · openclaw/openclaw@0b6c39b Auto-reply: avoid eager bundled route fallback · openclaw/openclaw@3ea1bf4 Tests: narrow session binding contract setup · openclaw/openclaw@54e4e16 fix(macOS): enable undo/redo in webchat composer text input (#34962) · openclaw/openclaw@00951dc Tests: speed up channel setup promotion · openclaw/openclaw@82b529a Docs: refresh agent instructions · openclaw/openclaw@5775fe2 fix(auth): serialize OAuth refresh across agents to fix #26322 (#67876) · openclaw/openclaw@8e79080 test: allow ollama public surface boundary test · openclaw/openclaw@7d4f1a6 Docs: add test performance guardrails · openclaw/openclaw@89706d3 Tests: restore context-engine usage proof · openclaw/openclaw@e4c4f95 Tests: slim context engine runtime coverage · openclaw/openclaw@74c198f ci: retry failed custom checkouts · openclaw/openclaw@0ee5baf test: trim duplicate provider auth onboarding cases · openclaw/openclaw@1ffc02e matrix: fix sessions_spawn --thread subagent session spawning (#67643) · openclaw/openclaw@1ce2596 test: reduce auth choice fixture churn · openclaw/openclaw@857b9cd test: mock health status config boundaries · openclaw/openclaw@9d5ab4a test: mock onboard config io boundary · openclaw/openclaw@299694d test: mock legacy state plugin boundaries · openclaw/openclaw@2713089 test: mock channel install boundaries · openclaw/openclaw@b945248 test: mock doctor preview channel boundaries · openclaw/openclaw@b1a3ad4
fix(outbound): centralize active delivery claims · opencl...
steipete · 2026-04-23 · via Recent Commits to openclaw:main

@@ -36,6 +36,10 @@ export interface PendingDeliveryDrainDecision {

3636

bypassBackoff?: boolean;

3737

}

383839+

export type ActiveDeliveryClaimResult<T> =

40+

| { status: "claimed"; value: T }

41+

| { status: "claimed-by-other-owner" };

42+3943

const MAX_RETRIES = 5;

40444145

/** Backoff delays in milliseconds indexed by retry count (1-based). */

@@ -90,20 +94,19 @@ function releaseRecoveryEntry(entryId: string): void {

9094

entriesInProgress.delete(entryId);

9195

}

929693-

/**

94-

* Claim an entry id against the shared in-memory recovery set so a concurrent

95-

* reconnect/startup drain will skip it while the owning caller is mid-flight.

96-

* Returns `false` if the id is already claimed. Callers must pair a successful

97-

* claim with {@link releaseActiveDelivery} in a `finally`. The claim is

98-

* process-local and intentionally does not survive a crash, so crash-replay

99-

* paths still recover fresh entries whose owning process died.

100-

*/

101-

export function tryClaimActiveDelivery(entryId: string): boolean {

102-

return claimRecoveryEntry(entryId);

103-

}

97+

export async function withActiveDeliveryClaim<T>(

98+

entryId: string,

99+

fn: () => Promise<T>,

100+

): Promise<ActiveDeliveryClaimResult<T>> {

101+

if (!claimRecoveryEntry(entryId)) {

102+

return { status: "claimed-by-other-owner" };

103+

}

104104105-

export function releaseActiveDelivery(entryId: string): void {

106-

releaseRecoveryEntry(entryId);

105+

try {

106+

return { status: "claimed", value: await fn() };

107+

} finally {

108+

releaseRecoveryEntry(entryId);

109+

}

107110

}

108111109112

function buildRecoveryDeliverParams(entry: QueuedDelivery, cfg: OpenClawConfig) {

@@ -246,15 +249,8 @@ export async function drainPendingDeliveries(opts: {

246249

const now = Date.now();

247250

const deliver = opts.deliver;

248251

const matchingEntries = (await loadPendingDeliveries(opts.stateDir))

249-

.map((entry) => ({

250-

entry,

251-

decision: opts.selectEntry(entry, now),

252-

}))

253-

.filter(

254-

(item): item is { entry: QueuedDelivery; decision: PendingDeliveryDrainDecision } =>

255-

item.decision.match,

256-

)

257-

.toSorted((a, b) => a.entry.enqueuedAt - b.entry.enqueuedAt);

252+

.filter((entry) => opts.selectEntry(entry, now).match)

253+

.toSorted((a, b) => a.enqueuedAt - b.enqueuedAt);

258254259255

if (matchingEntries.length === 0) {

260256

return;

@@ -264,7 +260,7 @@ export async function drainPendingDeliveries(opts: {

264260

`${opts.logLabel}: ${matchingEntries.length} pending message(s) matched ${opts.drainKey}`,

265261

);

266262267-

for (const { entry, decision } of matchingEntries) {

263+

for (const entry of matchingEntries) {

268264

if (!claimRecoveryEntry(entry.id)) {

269265

opts.log.info(`${opts.logLabel}: entry ${entry.id} is already being recovered`);

270266

continue;

@@ -280,6 +276,12 @@ export async function drainPendingDeliveries(opts: {

280276

continue;

281277

}

282278279+

const currentDecision = opts.selectEntry(currentEntry, Date.now());

280+

if (!currentDecision.match) {

281+

opts.log.info(`${opts.logLabel}: entry ${currentEntry.id} no longer matches, skipping`);

282+

continue;

283+

}

284+283285

if (currentEntry.retryCount >= MAX_RETRIES) {

284286

try {

285287

await moveToFailed(currentEntry.id, opts.stateDir);

@@ -296,7 +298,7 @@ export async function drainPendingDeliveries(opts: {

296298

continue;

297299

}

298300299-

if (!decision.bypassBackoff) {

301+

if (!currentDecision.bypassBackoff) {

300302

const retryEligibility = isEntryEligibleForRecoveryRetry(currentEntry, Date.now());

301303

if (!retryEligibility.eligible) {

302304

opts.log.info(