| 1 | import test from 'node:test'; |
| 2 | import assert from 'node:assert/strict'; |
| 3 | import fs from 'node:fs/promises'; |
| 4 | import path from 'node:path'; |
| 5 | import os from 'node:os'; |
| 6 | import vm from 'node:vm'; |
| 7 | import {ThreadStore} from '../src/lib.mjs'; |
| 8 | |
| 9 | // Execute the real startup/admission functions with local storage and inert |
| 10 | // runtime/delivery seams. No SDK boot, credentials, network, or model calls. |
| 11 | function extract(source,name) { |
| 12 | const start=source.search(new RegExp(`^(?:async )?function ${name}\\(`,'m')); |
| 13 | assert.notEqual(start,-1,`missing ${name}`); |
| 14 | const rest=source.slice(start),end=rest.slice(1).search(/\n(?:async )?function \w+\(/); |
| 15 | return end<0?rest:rest.slice(0,end+1); |
| 16 | } |
| 17 | for(const platform of ['telegram','feishu']) { |
| 18 | const lib=await import(`../../${platform}-bridge/src/lib.mjs`); |
| 19 | const source=await fs.readFile(new URL(`../../${platform}-bridge/src/index.mjs`,import.meta.url),'utf8'); |
| 20 | const direct=platform==='telegram'?'private':'p2p'; |
| 21 | const baseIdentity={chatId:'chat',chatType:direct,userId:'operator',username:'@operator',openId:'open-operator',unionId:'union-operator',isBot:false}; |
| 22 | async function fixture(t,identity=baseIdentity,policy={}) { |
| 23 | const dir=await fs.mkdtemp(path.join(os.tmpdir(),'bridge-recovery-'));t.after(()=>fs.rm(dir,{recursive:true,force:true})); |
| 24 | const file=path.join(dir,'thread-map.json'); |
| 25 | const writer=await ThreadStore.open(file,{messageLimit:200}); |
| 26 | await writer.setChat('chat',{threadId:'thread',activeTurnId:'turn',lastSeq:7,authorizedIdentity:identity,replyToMessageId:'original'}); |
| 27 | const threadStore=await ThreadStore.open(file,{messageLimit:200}); |
| 28 | const calls={runtime:[],sent:[],stream:[],commands:[]}; |
| 29 | const context=vm.createContext({...lib,threadStore,config:{allowlist:['operator'],allowGroups:false,allowUnlisted:false,requirePrefixInGroup:false,groupPrefix:'/cw',...policy}, |
| 30 | runtimeJson:async route=>{calls.runtime.push(route);return {turns:[{id:'turn',status:'in_progress'}]};}, |
| 31 | sendText:async(...args)=>calls.sent.push(args),sendTurnText:async(...args)=>calls.sent.push(args), |
| 32 | streamTurnEvents:async(...args)=>calls.stream.push(args),startTrackedTurnStream:(...args)=>calls.stream.push(args), |
| 33 | handleCommand:async(...args)=>calls.commands.push(args),answerCallback:async()=>{},callbackAction:()=>({kind:'status'}),handleModalAction:async(...args)=>calls.commands.push(args), |
| 34 | }); |
| 35 | const names=platform==='telegram'?['reattachActiveTurns','handleIncomingUpdate','handleCallbackQuery','rememberAuthorizedIdentity']:['reattachActiveTurns','handleIncomingMessage']; |
| 36 | for(const name of names) { |
| 37 | if(name==='rememberAuthorizedIdentity' && !source.includes('function rememberAuthorizedIdentity(')) continue; |
| 38 | vm.runInContext(extract(source,name),context); |
| 39 | } |
| 40 | return {context,calls,threadStore}; |
| 41 | } |
| 42 | test(`${platform}: restart rechecks current sender, group policy and saved provenance before any runtime read`,async t=>{ |
| 43 | for(const [identity,policy] of [ |
| 44 | [baseIdentity,{allowlist:['someone-else']}], |
| 45 | [{...baseIdentity,chatType:'group'},{}], |
| 46 | [null,{allowlist:['chat']}], |
| 47 | [null,{allowUnlisted:true}], |
| 48 | [{...baseIdentity,chatId:'other'},{}], |
| 49 | [{...baseIdentity,chatType:''},{}], |
| 50 | ...(platform==='telegram'?[[{...baseIdentity,isBot:true},{}]]:[]), |
| 51 | ]) { |
| 52 | const f=await fixture(t,identity,policy);await f.context.reattachActiveTurns(); |
| 53 | assert.deepEqual(f.calls,{runtime:[],sent:[],stream:[],commands:[]}); |
| 54 | } |
| 55 | for(const policy of [{allowlist:['operator']},{allowlist:['chat']},{allowUnlisted:true},{allowGroups:true}]) { |
| 56 | const identity=policy.allowGroups?{...baseIdentity,chatType:'group'}:baseIdentity; |
| 57 | const f=await fixture(t,identity,policy);await f.context.reattachActiveTurns(); |
| 58 | assert.deepEqual(f.calls.runtime,['/v1/threads/thread']);assert.equal(f.calls.sent.length,1);assert.equal(f.calls.stream.length,1); |
| 59 | } |
| 60 | }); |
| 61 | function incoming(userId,chatType=direct) { |
| 62 | return platform==='telegram' |
| 63 | ?{message:{message_id:1,chat:{id:'chat',type:chatType},from:{id:userId,username:'operator'},text:'/status'}} |
| 64 | :{sender:{sender_id:{user_id:userId}},message:{chat_id:'chat',chat_type:chatType,message_id:'new-reply',message_type:'text',content:JSON.stringify({text:'/status'})}}; |
| 65 | } |
| 66 | test(`${platform}: only admitted messages can update recovery and reply provenance`,async t=>{ |
| 67 | const f=await fixture(t); |
| 68 | const handle=platform==='telegram'?f.context.handleIncomingUpdate:f.context.handleIncomingMessage; |
| 69 | await handle(incoming('revoked')); |
| 70 | let state=await f.threadStore.getChat('chat'); |
| 71 | assert.equal(state.authorizedIdentity.userId,'operator');assert.equal(state.replyToMessageId,'original');assert.equal(f.calls.commands.length,0); |
| 72 | // Use a new ID: the first event was recorded as handled, as in production. |
| 73 | const allowed=incoming('operator');if(platform==='telegram')allowed.message.message_id=2;else allowed.message.message_id='allowed-reply'; |
| 74 | await handle(allowed);state=await f.threadStore.getChat('chat'); |
| 75 | assert.equal(state.authorizedIdentity.chatId,'chat');assert.equal(state.authorizedIdentity.userId,'operator');assert.equal(f.calls.commands.length,1); |
| 76 | assert.equal(lib.preservedChatStateFields(state).authorizedIdentity,state.authorizedIdentity,'thread replacement must retain authorization provenance'); |
| 77 | assert.equal(Object.hasOwn(state.authorizedIdentity,'text'),false); |
| 78 | if(platform==='feishu') assert.equal(state.replyToMessageId,'allowed-reply'); |
| 79 | }); |
| 80 | if(platform==='telegram') test('Telegram callbacks persist admitted identity for later recovery',async t=>{ |
| 81 | const f=await fixture(t,null); |
| 82 | await f.context.handleCallbackQuery({id:'callback',data:'status',message:{chat:{id:'chat',type:'private'},message_id:1},from:{id:'operator'}}); |
| 83 | assert.equal((await f.threadStore.getChat('chat')).authorizedIdentity.userId,'operator');assert.equal(f.calls.commands.length,1); |
| 84 | }); |
| 85 | } |
| 86 |