@@ -2,6 +2,11 @@ import { validateDecisions, applyDecisions } from "./dream/decisions.js";
22import { clusterMemories , findPotentialConflicts , cosineSimilarity } from "./dream/clustering.js" ;
33import { clusterByTag , intersectEvidence } from "./dream/narratives.js" ;
44import { scopeKeyOf } from "./scope.js" ;
5+ // Issue #239(第 4 项)镜像到巩固:错峰时段解析与「最近的高峰结束时刻」直接复用
6+ // 蒸馏侧已导出的纯函数,不另写一份解析器——两份实现漂移会让「同一个时段串在两处
7+ // 行为不同」,那比没有这个功能更糟。summarize.js 只依赖 dsh-llm 与 lang.js,
8+ // 不反向依赖 dream.js,无循环引用。
9+ import { isInPeakWindow , nextOffPeakAt } from "./summarize.js" ;
510import { createHash , randomUUID } from "node:crypto" ;
611import { STR , langOf } from "./lang.js" ;
712export { validateDecisions , applyDecisions , withEffortFallback , describeStreamFailure , resolveDreamEffort , resolveRoute } ;
@@ -744,8 +749,12 @@ export async function maintainIndexAfterDream(decisions, service, semantic) {
744749 if ( embedder . modelHash ) vectorIndex . markModel ?. ( embedder . modelHash , embedder . dimension ) ;
745750}
746751
747- export function createDreamScheduler ( { onRun, thresholdCount = 10 , thresholdChars = 5000 , delayMs = 2000 , minIntervalMs = 0 , logger, semantic = null , lastRunAtSeed = 0 } ) {
752+ export function createDreamScheduler ( { onRun, thresholdCount = 10 , thresholdChars = 5000 , delayMs = 2000 , minIntervalMs = 0 , logger, semantic = null , lastRunAtSeed = 0 , peakHours = "" , peakMaxDeferMinutes = 120 , auditPeakSkip = null , now = ( ) => Date . now ( ) , setTimeoutFn = setTimeout , clearTimeoutFn = clearTimeout } ) {
748753 let pendingTimer = null ;
754+ // Issue #239(第 4 项)镜像到巩固:高峰顺延定时器。与 pendingTimer 分开——两者
755+ // 语义不同(一个是「马上要跑」,一个是「等出高峰再跑」),合成一个变量会让
756+ // maybeSchedule 的守卫在顺延期间把新的写入触发误当成「已有待跑」而吞掉。
757+ let deferTimer = null ;
749758 let running = false ;
750759 let disposed = false ;
751760 let baseline = { count : 0 , chars : 0 } ;
@@ -766,51 +775,103 @@ export function createDreamScheduler({ onRun, thresholdCount = 10, thresholdChar
766775 }
767776
768777 function maybeSchedule ( service ) {
769- if ( disposed || running || pendingTimer ) return false ;
778+ if ( disposed || running || pendingTimer || deferTimer ) return false ;
770779 // Issue #89(请求 2):最小触发间隔闸门。
771- if ( minIntervalMs > 0 && Date . now ( ) - lastRunAt < minIntervalMs ) return false ;
780+ if ( minIntervalMs > 0 && now ( ) - lastRunAt < minIntervalMs ) return false ;
772781 const { trigger, count, chars } = shouldTrigger ( service ) ;
773782 if ( ! trigger ) return false ;
774- pendingTimer = setTimeout ( ( ) => {
783+ // Issue #239(第 4 项,错峰队列)镜像到巩固:命中高峰就不调模型。与蒸馏的差别
784+ // 在于巩固是**全局单实例**(蒸馏按会话各挂一个定时器),所以这里只需要一个
785+ // deferTimer,且不需要 deferredRuns 那套按会话去重。
786+ // baseline 刻意不刷新:阈值继续累积,留到非高峰一次性巩固(一次大 run 比多次
787+ // 小 run 省)。审计只登记一行 skip——「为什么不再做梦了」必须对用户可观测。
788+ if ( isInPeakWindow ( new Date ( now ( ) ) , peakHours ) ) {
789+ try {
790+ auditPeakSkip ?. ( { count, chars } ) ;
791+ } catch ( error ) {
792+ // 记账是 best-effort:写审计行失败只 warn,绝不反噬调度本身。
793+ logger ?. warn ?. ( `dsh-mneme dream: peak-hours audit failed: ${ String ( error ) } ` ) ;
794+ }
795+ scheduleDeferredRun ( service ) ;
796+ return false ;
797+ }
798+ pendingTimer = setTimeoutFn ( ( ) => {
775799 pendingTimer = null ;
776- running = true ;
777- lastRunAt = Date . now ( ) ;
778- // Defer the onRun invocation so a synchronous throw cannot escape the
779- // timer callback (which would crash the process) and skip the teardown.
780- // Errors are logged, never swallowed silently. inFlight lets dispose()
781- // await the running consolidation before the caller closes the store.
782- inFlight = Promise . resolve ( )
783- . then ( ( ) => ( onRun ? onRun ( ) : Promise . resolve ( { ok : true , skipped : true } ) ) )
784- . then ( ( result ) => {
785- // Refresh the baseline only for a successful run (design §5.3: an
786- // LLM failure must not move the baseline, so the next write can
787- // immediately re-trigger a retry). A `{ok:false}` result or a throw
788- // keeps the old baseline. A run that reports nothing is treated as
789- // completed without failure (no-op hooks / minimal test doubles).
790- if ( result && result . ok ) {
791- try {
792- baseline = shouldTrigger ( service ) ;
793- } catch ( error ) {
794- // Store closed mid-flight: keep the last known baseline.
795- logger ?. warn ?. ( `dsh-mneme dream: baseline refresh failed: ${ String ( error ) } ` ) ;
796- }
797- }
798- } )
799- . catch ( ( error ) => {
800- logger ?. warn ?. ( `dsh-mneme dream: run failed: ${ error ?. message ?? error } ` ) ;
801- // Failed runs do not refresh the baseline.
802- } )
803- . finally ( ( ) => {
804- running = false ;
805- inFlight = null ;
806- } ) ;
800+ startRun ( service ) ;
807801 } , delayMs ) ;
808802 return true ;
809803 }
810804
805+ /**
806+ * Issue #239:高峰内择时补跑——挂到「距当前最近的一个高峰结束时刻」,被
807+ * peakMaxDeferMinutes 截断时到点照跑(bypassPeak),长高峰不会把巩固饿死。
808+ * 定时器 unref:不阻止宿主退出。重复触发不叠加(deferTimer 已在 maybeSchedule
809+ * 的守卫里,这里再判一次以防从其它路径进来)。
810+ */
811+ function scheduleDeferredRun ( service ) {
812+ if ( disposed || deferTimer ) return ;
813+ const at = nextOffPeakAt ( new Date ( now ( ) ) , peakHours ) ;
814+ if ( ! at ) return ;
815+ const maxDeferMs = ( peakMaxDeferMinutes ?? 0 ) * 60000 ;
816+ let delay = Math . max ( 0 , at . getTime ( ) - now ( ) ) ;
817+ const capped = maxDeferMs > 0 && delay > maxDeferMs ;
818+ if ( capped ) delay = maxDeferMs ;
819+ deferTimer = setTimeoutFn ( ( ) => {
820+ deferTimer = null ;
821+ if ( disposed ) return ;
822+ // 截断放行时仍在高峰:不再重新顺延(否则长高峰里会无限顺延,等于把巩固
823+ // 关掉)。直接开跑,与蒸馏的 bypassPeak 同口径。
824+ if ( ! capped && isInPeakWindow ( new Date ( now ( ) ) , peakHours ) ) {
825+ // 理论上到点已出高峰;时钟跳变/时段串被改小可能落回高峰内,此时再顺延一次。
826+ scheduleDeferredRun ( service ) ;
827+ return ;
828+ }
829+ logger ?. info ?. ( `dsh-mneme dream: peak-hours deferred run firing (capped=${ capped } , delayMs=${ delay } )` ) ;
830+ startRun ( service ) ;
831+ } , delay ) ;
832+ deferTimer . unref ?. ( ) ;
833+ }
834+
835+ /**
836+ * 真正开跑。抽出来是因为两条路径都要用:写入触发的正常路径,与高峰顺延后的
837+ * 补跑路径。onRun 的调用刻意放在 Promise 里——同步抛出的异常若逃出 timer 回调
838+ * 会直接崩掉进程并跳过收尾。inFlight 让 dispose() 能等完这一轮再关库。
839+ */
840+ function startRun ( service ) {
841+ running = true ;
842+ lastRunAt = now ( ) ;
843+ inFlight = Promise . resolve ( )
844+ . then ( ( ) => ( onRun ? onRun ( ) : Promise . resolve ( { ok : true , skipped : true } ) ) )
845+ . then ( ( result ) => {
846+ // Refresh the baseline only for a successful run (design §5.3: an
847+ // LLM failure must not move the baseline, so the next write can
848+ // immediately re-trigger a retry). A `{ok:false}` result or a throw
849+ // keeps the old baseline. A run that reports nothing is treated as
850+ // completed without failure (no-op hooks / minimal test doubles).
851+ if ( result && result . ok ) {
852+ try {
853+ baseline = shouldTrigger ( service ) ;
854+ } catch ( error ) {
855+ // Store closed mid-flight: keep the last known baseline.
856+ logger ?. warn ?. ( `dsh-mneme dream: baseline refresh failed: ${ String ( error ) } ` ) ;
857+ }
858+ }
859+ } )
860+ . catch ( ( error ) => {
861+ logger ?. warn ?. ( `dsh-mneme dream: run failed: ${ error ?. message ?? error } ` ) ;
862+ // Failed runs do not refresh the baseline.
863+ } )
864+ . finally ( ( ) => {
865+ running = false ;
866+ inFlight = null ;
867+ } ) ;
868+ }
869+
811870 async function dispose ( ) {
812871 disposed = true ;
813- if ( pendingTimer ) { clearTimeout ( pendingTimer ) ; pendingTimer = null ; }
872+ if ( pendingTimer ) { clearTimeoutFn ( pendingTimer ) ; pendingTimer = null ; }
873+ // Issue #239:高峰顺延定时器同样要清,否则进程关闭后仍会触发一次巩固。
874+ if ( deferTimer ) { clearTimeoutFn ( deferTimer ) ; deferTimer = null ; }
814875 // An in-flight run is left to complete naturally (its LLM calls are
815876 // already paid for and aborting would discard the work). Await it so the
816877 // caller can close the store only after every write has landed.
0 commit comments