package vip.mate.channel.wecom;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Component;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.ScheduledFuture;
import java.util.concurrent.TimeUnit;
/**
* Periodically refreshes a WeCom AI Bot {@code stream} reply with the
* "π€ ζθδΈ..." placeholder text so WeCom's server-side does not drop
* the stream slot while a long-running agent task is still computing.
*
*
Why this exists: WeCom's stream slot has an undocumented TTL
* (empirically observed at ~60-120s of silence). When the slot drops,
* the eventual {@code finish=true} chunk is silently rejected β the
* user sees "π€ ζθδΈ..." stuck forever. RFC-32 Β§2.1.2 (R-7 / B-5).
*
*
Constants (chosen empirically based on observed slot lifetime):
*
* - 20s refresh interval β well under the observed 60s minimum drop
* - 180s force-finish ceiling β bound the worst-case "stuck" UX even
* if the agent task hangs forever; lets the next user reply use a
* fresh stream slot rather than reuse a closed one
*
*
* Force-finish invariant: after the 180s ceiling fires, the
* scheduler calls {@link WeComChannelAdapter#invalidateReplyContext}
* to evict the {@code (frameReqId, processingStreamId)} pair from
* {@code replyContexts} β so when the eventual real reply arrives,
* {@code renderAndSend} sees no context and falls through to a fresh
* {@code sendMessage} path instead of reusing the dead stream id.
* RFC-32 Β§2.1.2 / R-7 closes this race; the adapter's
* {@code invalidateReplyContext} is the contract.
*/
@Slf4j
@Component
public class WeComKeepaliveScheduler {
/** Refresh interval (seconds). */
static final long REFRESH_INTERVAL_SECONDS = 20;
/** Hard ceiling β after this many seconds, force-finish the stream. */
static final long MAX_DURATION_SECONDS = 180;
/** Placeholder text written on every refresh tick + on force-finish. */
static final String PROCESSING_TEXT = "π€ ζθδΈ...";
/** One-shot state per active stream. Held by reference inside the scheduled task. */
private static final class StreamState {
final WeComChannelAdapter adapter;
final String reqId;
final String streamId;
final String replyToken;
final long startedAt;
volatile ScheduledFuture> future;
StreamState(WeComChannelAdapter a, String r, String s, String t) {
this.adapter = a; this.reqId = r; this.streamId = s; this.replyToken = t;
this.startedAt = System.currentTimeMillis();
}
}
private final Map states = new ConcurrentHashMap<>();
private final ScheduledExecutorService scheduler = Executors.newScheduledThreadPool(2, r -> {
Thread t = new Thread(r, "wecom-keepalive");
t.setDaemon(true);
return t;
});
/**
* Begin keepalive for a freshly-issued processing stream. No-op if
* called twice for the same {@code streamId} (just logs and keeps
* the existing schedule alive).
*
* @param adapter the live WeCom adapter (carries replyStream + invalidateReplyContext)
* @param reqId the inbound message frame's req_id (binds the
* outbound reply_stream chunks)
* @param streamId the same stream_id used in the initial "π€ ζθδΈ..." chunk
* @param replyToken target id used to look up reply context on
* invalidation (typically chatId for groups,
* senderId for direct messages β same value passed
* to {@code replyContexts.put})
*/
public void start(WeComChannelAdapter adapter, String reqId, String streamId, String replyToken) {
if (adapter == null || reqId == null || streamId == null
|| reqId.isBlank() || streamId.isBlank()) {
log.debug("[wecom-keepalive] start ignored β null/blank arg(s)");
return;
}
if (states.containsKey(streamId)) {
log.debug("[wecom-keepalive] start ignored β stream {} already tracked", streamId);
return;
}
StreamState st = new StreamState(adapter, reqId, streamId, replyToken);
st.future = scheduler.scheduleAtFixedRate(
() -> tick(st),
REFRESH_INTERVAL_SECONDS,
REFRESH_INTERVAL_SECONDS,
TimeUnit.SECONDS);
states.put(streamId, st);
log.debug("[wecom-keepalive] started for stream={} reqId={}", streamId, reqId);
}
/**
* Stop keepalive for a stream β call this immediately before sending
* the real reply so the next refresh tick doesn't race the
* {@code finish=true} chunk on the same stream.
*/
public void stop(String streamId) {
if (streamId == null || streamId.isBlank()) return;
StreamState st = states.remove(streamId);
if (st != null && st.future != null) {
st.future.cancel(false);
log.debug("[wecom-keepalive] stopped for stream={}", streamId);
}
}
/**
* Drop every tracked stream and cancel its schedule. Called from
* {@code releaseConnectionResources} so reconnects start clean.
*/
public void shutdownAll() {
for (StreamState st : states.values()) {
if (st.future != null) st.future.cancel(false);
}
states.clear();
}
private void tick(StreamState st) {
long elapsedSec = (System.currentTimeMillis() - st.startedAt) / 1000;
if (elapsedSec >= MAX_DURATION_SECONDS) {
// Force-finish: send finish=true on the same stream so the
// server-side closes the slot cleanly, then evict the
// replyContext entry so the eventual real reply takes the
// fresh-stream path.
try {
st.adapter.replyStreamFinishForKeepalive(st.reqId, st.streamId, PROCESSING_TEXT);
} catch (Exception e) {
log.debug("[wecom-keepalive] force-finish replyStream failed for {}: {}",
st.streamId, e.getMessage());
}
try {
st.adapter.invalidateReplyContext(st.replyToken, st.streamId);
} catch (Exception e) {
log.debug("[wecom-keepalive] invalidateReplyContext failed for {}: {}",
st.streamId, e.getMessage());
}
stop(st.streamId);
log.info("[wecom-keepalive] force-finished stream {} after {}s ceiling",
st.streamId, MAX_DURATION_SECONDS);
return;
}
try {
st.adapter.replyStreamRefreshForKeepalive(st.reqId, st.streamId, PROCESSING_TEXT);
} catch (Exception e) {
log.debug("[wecom-keepalive] refresh failed for {}: {}", st.streamId, e.getMessage());
}
}
// ---- Test hooks ----
int activeStreamCount() { return states.size(); }
}