diff --git a/aesel/docs/jev-harness.md b/aesel/docs/jev-harness.md index ad056b5260..babbcfa080 100644 --- a/aesel/docs/jev-harness.md +++ b/aesel/docs/jev-harness.md @@ -141,3 +141,31 @@ braincells. Phone receipts are saved only during explicit debug fixture runs. Live results and failed attempts are retained under `apple/walkieware/Tests/audio/`. Decision latency is not an end-to-end creation speed claim: the existing coding model still generates the piece. + +### Persistent phone input + +Lith's existing DigitalOcean process serves `wss://aesthetic.computer/api/easel-musical-stream`. +The phone opens it after sign-in and closes it while backgrounded. Authentication +is the first JSON message (`{type:"authenticate",token}`), never a URL query; +`ready` confirms the session. Sessions renew after ten minutes. A foreground +client reconnects after disconnects, and HTTP remains available while cold. +Once a decision has been sent, a disconnect does not replay it over HTTP. + +`{type:"observation",body:}` receives an immediate +`received` acknowledgement followed by a `decision` or `error`, both with the +original session and sequence. One decision runs per socket; only the latest +waiting observation survives. `{type:"cancel",sessionId}` aborts active work +and drops pending work for that capture. Heartbeats, bounded messages, two +connections per account and a global connection cap bound the transport. +Authentication is reused; durable per-decision quotas and the shared account +concurrency guard still apply across HTTP and sockets. + +Receipts separate acknowledgement round-trip (`transportMs`), server processing +including quota checks (`serverMs`), and provider request duration (`providerMs`). +Acknowledgement time includes scheduling and is not a pure network measurement. +The upstream Decisions API remains HTTP. Neither this transport nor a separate +VM removes provider inference latency. Measure these components before moving +regions or provisioning another service. + +Focused validation: `node --test lith/musical-socket.test.mjs aesel/test/musical-jev.test.mjs`. +Phone validation: `node apple/walkieware/Tests/musical-benchmark.mjs mixed-request --jev --socket`. diff --git a/aesel/src/musical-input-advisor.mjs b/aesel/src/musical-input-advisor.mjs index 844d7b2fb7..25b5855d22 100644 --- a/aesel/src/musical-input-advisor.mjs +++ b/aesel/src/musical-input-advisor.mjs @@ -20,8 +20,8 @@ export class MusicalInputAdvisor { const deadline=new Promise((_,reject)=>{timer=setTimeout(()=>{controller.abort();reject(Error('deadline'));},this.timeoutMs);}); const response=await Promise.race([deadline,this.fetchImpl(this.endpoint,{method:'POST',signal:controller.signal,headers:{'Content-Type':'application/json',Authorization:`Bearer ${bearer}`},body:JSON.stringify({schema:'walkieware-input/v1',sessionId,sequence,features})}).then(async r=>{if(!r.ok)throw Error(`http_${r.status}`);return r.json();})]); if(this.sessionId!==sessionId||controller.signal.aborted||response.schema!=='walkieware-decision/v1'||response.sessionId!==sessionId||response.sequence!==sequence||!Object.hasOwn(MUSICAL_CHOICES,response.choice)||!Number.isFinite(response.confidence)||response.confidence<.8||response.confidence>1)return null; - const value={key,choice:response.choice,cue:MUSICAL_CHOICES[response.choice],elapsedMs:Math.round(performance.now()-started)}; - this.cache=value;this.onEvent('jevDecision',{choice:value.choice,elapsedMs:value.elapsedMs,sequence});return value; + const value={key,choice:response.choice,cue:MUSICAL_CHOICES[response.choice],elapsedMs:Math.round(performance.now()-started),transport:response.transport||'http',transportMs:response.transportMs,serverMs:response.serverMs,providerMs:response.elapsedMs}; + this.cache=value;this.onEvent('jevDecision',{choice:value.choice,elapsedMs:value.elapsedMs,sequence,transport:response.transport||'http',transportMs:response.transportMs,serverMs:response.serverMs,providerMs:response.elapsedMs});return value; }catch(error){if(this.sessionId===sessionId)this.onEvent('jevFallback',{reason:/^http_\d+$/.test(error.message)?error.message:controller.signal.aborted?'timeout_or_cancel':error.name||'unavailable'});return null;} finally{clearTimeout(timer);if(this.sessionId===sessionId)this.pending=null;} })(); diff --git a/aesel/src/musical-input-socket.mjs b/aesel/src/musical-input-socket.mjs new file mode 100644 index 0000000000..a3232690e0 --- /dev/null +++ b/aesel/src/musical-input-socket.mjs @@ -0,0 +1,48 @@ +// One foreground connection; existing HTTP path remains the cold/offline fallback. +export class MusicalInputSocket { + constructor({token,onEvent=()=>{},WebSocketImpl=globalThis.WebSocket,url='wss://aesthetic.computer/api/easel-musical-stream',fetchImpl=(...args)=>globalThis.fetch(...args)}={}) { + Object.assign(this,{token,onEvent,WebSocketImpl,url,fetchImpl});this.pending=new Map();this.ready=false;this.enabled=true; + } + connect(){ + const bearer=this.token?.(); + if(!this.enabled||!bearer||!this.WebSocketImpl)return; + if(this.ws&&this.bearer===bearer)return; + this.close();this.bearer=bearer; + let ws;try{ws=new this.WebSocketImpl(this.url);}catch{return;} + this.ws=ws;const started=performance.now(); + const timer=setTimeout(()=>{if(this.ws===ws&&!this.ready)ws.close();},4000); + ws.onopen=()=>ws.send(JSON.stringify({type:'authenticate',token:bearer})); + ws.onmessage=event=>{ + if(this.ws!==ws)return; + let m;try{m=JSON.parse(event.data);}catch{return;} + if(m.type==='ready'){clearTimeout(timer);this.ready=true;this.onEvent('inputSocketReady',{connectMs:Math.round(performance.now()-started)});return;} + const key=`${m.sessionId}:${m.sequence}`,p=this.pending.get(key);if(!p)return; + if(m.type==='received'){p.transportMs=Math.round(performance.now()-p.started);this.onEvent('inputSocketAck',{transportMs:p.transportMs});return;} + if(m.type==='decision'){this.pending.delete(key);p.cleanup();p.resolve({ok:true,json:async()=>({...m,transport:'socket',transportMs:p.transportMs})});} + else if(m.type==='error'){this.pending.delete(key);p.cleanup();p.resolve({ok:false,status:m.status});} + }; + ws.onerror=()=>{}; + ws.onclose=()=>{clearTimeout(timer);if(this.ws!==ws)return;this.ws=null;this.ready=false;this.rejectPending();if(this.enabled)this.reconnect=setTimeout(()=>this.connect(),2000);}; + } + rejectPending(){for(const p of this.pending.values()){p.cleanup();p.reject(Error('socket_closed'));}this.pending.clear();} + close(){clearTimeout(this.reconnect);const ws=this.ws;this.ws=null;this.ready=false;this.rejectPending();ws?.close();} + suspend(){this.enabled=false;this.close();} + resume(){this.enabled=true;this.connect();} + fetch=async(endpoint,options)=>{ + this.connect(); + if(!this.ready){this.onEvent('inputHttpFallback',{reason:'socket_not_ready'});return this.fetchImpl(endpoint,options);} + const ws=this.ws,body=JSON.parse(options.body),key=`${body.sessionId}:${body.sequence}`; + // Once sent, never replay a possibly billed decision over HTTP. + return new Promise((resolve,reject)=>{ + const cancel=()=>{ + this.pending.delete(key);options.signal?.removeEventListener('abort',cancel); + if(ws.readyState===1)ws.send(JSON.stringify({type:'cancel',sessionId:body.sessionId})); + reject(Error('cancelled')); + }; + if(options.signal?.aborted){reject(Error('cancelled'));return;} + const cleanup=()=>options.signal?.removeEventListener('abort',cancel); + this.pending.set(key,{resolve,reject,cleanup,started:performance.now()});options.signal?.addEventListener('abort',cancel,{once:true}); + try{ws.send(JSON.stringify({type:'observation',body}));}catch{this.pending.delete(key);cleanup();reject(Error('socket_send_failed'));} + }); + }; +} diff --git a/aesel/test/musical-jev.test.mjs b/aesel/test/musical-jev.test.mjs index c571a0ec13..b0e41ef732 100644 --- a/aesel/test/musical-jev.test.mjs +++ b/aesel/test/musical-jev.test.mjs @@ -38,3 +38,8 @@ test('default fetch keeps its global receiver for Safari',async t=>{ t.mock.method(globalThis,'fetch',async function(_,o){assert.equal(this,globalThis);const b=JSON.parse(o.body);return Response.json({...b,schema:'walkieware-decision/v1',choice:'follow_speech',confidence:.95});}); const advisor=new MusicalInputAdvisor({token:()=> 'test'});assert.equal((await advisor.finish(input)).choice,'follow_speech'); }); +test('changed final observations request fresh advice instead of reusing the earlier state',async()=>{ + const bodies=[];const a=new MusicalInputAdvisor({token:()=> 'test',fetchImpl:async(_,o)=>{const b=JSON.parse(o.body);bodies.push(b);return Response.json({...b,schema:'walkieware-decision/v1',choice:b.features.hasSpeech?'follow_speech':'trace_pitch',confidence:.95});}}); + a.observe(input);await a.pending.work;const result=await a.finish({...input,transcript:'',words:[]}); + assert.equal(result.choice,'trace_pitch');assert.equal(bodies.length,2);assert.notEqual(bodies[0].sessionId,bodies[1].sessionId); +}); diff --git a/lith/musical-socket.mjs b/lith/musical-socket.mjs new file mode 100644 index 0000000000..2d0ce84ea9 --- /dev/null +++ b/lith/musical-socket.mjs @@ -0,0 +1,67 @@ +import {WebSocketServer, WebSocket} from 'ws'; + +// Same decision service and durable allowance as HTTP. No audio or tokens in URLs. +export function attachMusicalSocket(server, {authenticate, decide, authMs=5000, lifetimeMs=600000}={}) { + const wss=new WebSocketServer({noServer:true,maxPayload:8192,perMessageDeflate:false}); + const accounts=new Map(); + const upgrade=(req,socket,head)=>{ + if(req.url!=='/api/easel-musical-stream'||wss.clients.size>=64){socket.end('HTTP/1.1 503 Service Unavailable\r\nConnection: close\r\n\r\n');return;} + wss.handleUpgrade(req,socket,head,ws=>wss.emit('connection',ws)); + }; + server.on('upgrade',upgrade); + wss.on('connection',ws=>{ + let subject,authenticating=false,active=null,queued=null,closed=false,alive=true; + let messages=0,windowAt=Date.now(); + const send=value=>{if(ws.readyState===WebSocket.OPEN){if(ws.bufferedAmount>16384)ws.close(1008,'Slow reader');else ws.send(JSON.stringify(value));}}; + const authTimer=setTimeout(()=>ws.close(1008,'Authenticate first'),authMs); + const lifeTimer=setTimeout(()=>ws.close(1000,'Renew session'),lifetimeMs); + const pingTimer=setInterval(()=>{if(!alive){ws.terminate();return;}alive=false;ws.ping();},20000); + for(const t of [authTimer,lifeTimer,pingTimer])t.unref?.(); + ws.on('pong',()=>{alive=true;}); + const fail=(body,status)=>send({type:'error',sessionId:body?.sessionId,sequence:body?.sequence,status}); + async function run(body){ + const task={body,controller:new AbortController()};active=task; + const started=performance.now(); + try{ + const r=await decide({httpMethod:'POST',headers:{authorization:'socket-session'},body:JSON.stringify(body)},{subject,signal:task.controller.signal}); + if(!closed&&!task.controller.signal.aborted){ + if(r.statusCode===200)send({...JSON.parse(r.body),type:'decision',serverMs:Math.round(performance.now()-started)}); + else fail(body,r.statusCode); + } + }catch{if(!closed&&!task.controller.signal.aborted)fail(body,503);} + finally{active=null;if(queued&&!closed){const next=queued;queued=null;void run(next);}} + } + ws.on('message',async(data,binary)=>{ + if(binary){ws.close(1003,'JSON only');return;} + if(Date.now()-windowAt>=1000){messages=0;windowAt=Date.now();} + if(++messages>30){ws.close(1008,'Input rate exceeded');return;} + let m;try{m=JSON.parse(data.toString());}catch{ws.close(1008,'Invalid message');return;} + if(!m||typeof m!=='object'){ws.close(1008,'Invalid message');return;} + if(!subject){ + if(authenticating||m.type!=='authenticate'||typeof m.token!=='string'||m.token.length>7000){ws.close(1008,'Authenticate first');return;} + authenticating=true; + try{ + const user=await authenticate({authorization:`Bearer ${m.token}`}); + if(closed)return; + if(!user||(accounts.get(user)||0)>=2){ws.close(1008,'Account unavailable');return;} + subject=user;accounts.set(subject,(accounts.get(subject)||0)+1);clearTimeout(authTimer);send({type:'ready'}); + }catch{ws.close(1011,'Account unavailable');} + return; + } + if(m.type==='ping'){send({type:'pong'});return;} + if(m.type==='cancel'){ + if(active?.body.sessionId===m.sessionId)active.controller.abort(); + if(queued?.sessionId===m.sessionId)queued=null; + return; + } + if(m.type!=='observation'||!m.body||JSON.stringify(m.body).length>2048){ws.close(1008,'Invalid observation');return;} + const body=m.body; + if(!body||typeof body.sessionId!=='string'||!Number.isInteger(body.sequence)){ws.close(1008,'Invalid observation');return;} + send({type:'received',sessionId:body.sessionId,sequence:body.sequence}); + if(active){if(queued)fail(queued,409);queued=body;}else void run(body); + }); + ws.on('error',()=>{}); + ws.on('close',()=>{closed=true;clearTimeout(authTimer);clearTimeout(lifeTimer);clearInterval(pingTimer);active?.controller.abort();queued=null;if(subject){const n=(accounts.get(subject)||1)-1;if(n)accounts.set(subject,n);else accounts.delete(subject);}}); + }); + return {close(){server.off('upgrade',upgrade);for(const ws of wss.clients)ws.terminate();wss.close();},wss}; +} diff --git a/lith/musical-socket.test.mjs b/lith/musical-socket.test.mjs new file mode 100644 index 0000000000..d8ec40174f --- /dev/null +++ b/lith/musical-socket.test.mjs @@ -0,0 +1,59 @@ +import test from 'node:test'; +import assert from 'node:assert/strict'; +import {createServer} from 'node:http'; +import {once} from 'node:events'; +import {WebSocket} from 'ws'; +import {attachMusicalSocket} from './musical-socket.mjs'; +import {createMusicalHandler} from '../system/backend/easel-musical-jev.mjs'; +import {MusicalInputSocket} from '../aesel/src/musical-input-socket.mjs'; +import {MusicalInputAdvisor} from '../aesel/src/musical-input-advisor.mjs'; +const input={transcript:'Private words',words:[],sound:{frames:[1,2,3].map(i=>({atMs:i*100,pitchHz:440,rms:.2})),onsetsMs:[]}}; +const features={hasSpeech:true,hasTonalSound:true,soundAfterSpeech:false,contour:'steady',attacks:0,rhythm:'unknown',energy:'steady'}; +const body=(sequence=1)=>({schema:'walkieware-input/v1',sessionId:'12345678-1234-1234-1234-123456789012',sequence,features}); +const answer={answers:{mapping:{choice:'follow_speech',probabilities:{follow_speech:.95}}}}; +async function fixture(t,overrides={}){ + let auths=0,calls=0,quotas=0; + const authenticate=async h=>{auths++;return h.authorization==='Bearer valid'?'account':null;}; + const decide=createMusicalHandler({authenticate,budget:{consume:async()=>{quotas++;return true;}},evaluate:async()=>{calls++;return answer;}}); + const server=createServer();const binding=attachMusicalSocket(server,{authenticate,decide,...overrides});server.listen(0,'127.0.0.1');await once(server,'listening'); + t.after(()=>{binding.close();server.close();}); + return {url:`ws://127.0.0.1:${server.address().port}/api/easel-musical-stream`,counts:()=>({auths,calls,quotas})}; +} +function inbox(ws){const messages=[],waiters=[];ws.on('message',d=>{const m=JSON.parse(d);const i=waiters.findIndex(w=>w.type===m.type);if(i>=0)waiters.splice(i,1)[0].resolve(m);else messages.push(m);});return type=>{const i=messages.findIndex(m=>m.type===type);if(i>=0)return Promise.resolve(messages.splice(i,1)[0]);return new Promise(resolve=>waiters.push({type,resolve}));};} +async function client(url){const ws=new WebSocket(url),next=inbox(ws);await once(ws,'open');return {ws,next,send:m=>ws.send(JSON.stringify(m))};} +test('real socket authenticates once, keeps per-decision quotas and reports timing', {timeout:3000},async t=>{ + const f=await fixture(t),c=await client(f.url);c.send({type:'authenticate',token:'valid'});await c.next('ready'); + for(let sequence=1;sequence<=2;sequence++){c.send({type:'observation',body:body(sequence)});assert.equal((await c.next('received')).sequence,sequence);const r=await c.next('decision');assert.equal(r.choice,'follow_speech');assert.equal(r.sequence,sequence);assert.ok(r.serverMs>=r.elapsedMs);} + assert.deepEqual(f.counts(),{auths:1,calls:2,quotas:2}); + c.send({type:'observation',body:{...body(3),prompt:'injected'}});assert.equal((await c.next('error')).status,400);assert.equal(f.counts().calls,2); +}); +test('unauthenticated, malformed and expired connections are closed', {timeout:3000},async t=>{ + const f=await fixture(t,{authMs:30,lifetimeMs:150}); + const a=await client(f.url);a.send({type:'observation',body:body()});assert.equal((await once(a.ws,'close'))[0],1008); + const b=await client(f.url);b.send({type:'authenticate',token:'wrong'});assert.equal((await once(b.ws,'close'))[0],1008); + const c=await client(f.url);assert.equal((await once(c.ws,'close'))[0],1008); + const d=await client(f.url);d.send({type:'authenticate',token:'valid'});await d.next('ready');d.send({type:'observation'});assert.equal((await once(d.ws,'close'))[0],1008); + const e=await client(f.url);e.send({type:'authenticate',token:'valid'});await e.next('ready');assert.equal((await once(e.ws,'close'))[0],1000); + assert.equal(f.counts().calls,0); +}); +test('slow inference keeps newest pending observation; cancel aborts active work', {timeout:3000},async t=>{ + const started=[],releases=[]; + const f=await fixture(t,{decide:(event,{signal})=>new Promise(resolve=>{const b=JSON.parse(event.body);started.push(b.sequence);const release=()=>resolve({statusCode:200,body:JSON.stringify({...b,schema:'walkieware-decision/v1',choice:'sustain',confidence:.95})});releases.push(release);signal.addEventListener('abort',release);})}); + const c=await client(f.url);c.send({type:'authenticate',token:'valid'});await c.next('ready'); + c.send({type:'observation',body:body(1)});await c.next('received'); + c.send({type:'observation',body:body(2)});await c.next('received'); + c.send({type:'observation',body:body(3)});await c.next('received');assert.equal((await c.next('error')).sequence,2); + releases[0]();assert.equal((await c.next('decision')).sequence,1);assert.deepEqual(started,[1,3]); + c.send({type:'cancel',sessionId:body().sessionId});await new Promise(r=>setTimeout(r,10)); + c.send({type:'observation',body:body(4)});await c.next('received');assert.deepEqual(started,[1,3,4]);releases[2]();assert.equal((await c.next('decision')).sequence,4); +}); +test('phone client uses warm socket, matching cache, HTTP on cold connection, reconnect after close', {timeout:4000},async t=>{ + const f=await fixture(t);const events=[];let ready;const waitReady=()=>new Promise(r=>ready=r);let readyPromise=waitReady(); + const transport=new MusicalInputSocket({url:f.url,WebSocketImpl:WebSocket,token:()=> 'valid',onEvent:(event,fields)=>{events.push({event,fields});if(event==='inputSocketReady')ready();},fetchImpl:async()=>{throw Error('Unexpected HTTP');}});t.after(()=>transport.suspend()); + transport.connect();await readyPromise; + const a=new MusicalInputAdvisor({token:()=> 'valid',fetchImpl:transport.fetch,onEvent:(event,fields)=>events.push({event,fields})}); + assert.equal((await a.finish(input)).choice,'follow_speech');assert.equal((await a.finish(input)).choice,'follow_speech'); + assert.equal(f.counts().calls,1);assert.equal(events.find(e=>e.event==='jevDecision').fields.transport,'socket');assert.ok(events.some(e=>e.event==='inputSocketAck')); + readyPromise=waitReady();transport.ws.close();await readyPromise;assert.equal(f.counts().auths,2); + let http=0;const cold=new MusicalInputSocket({token:()=> 'valid',WebSocketImpl:null,fetchImpl:async()=>{http++;return {ok:true};}});assert.equal((await cold.fetch('url',{body:'{}'})).ok,true);assert.equal(http,1);cold.suspend(); +}); diff --git a/lith/package-lock.json b/lith/package-lock.json index a963d382e3..a117705461 100644 --- a/lith/package-lock.json +++ b/lith/package-lock.json @@ -10,7 +10,8 @@ "dependencies": { "express": "^5.1.0", "mailparser": "^3.7.0", - "smtp-server": "^3.13.0" + "smtp-server": "^3.13.0", + "ws": "^8.21.3" }, "engines": { "node": ">=24.18.1 <25" @@ -1167,6 +1168,27 @@ "resolved": "https://registry.npmjs.org/wrappy/-/wrappy-1.0.2.tgz", "integrity": "sha512-l4Sp/DRseor9wL6EvV2+TuQn63dMkPjZ/sp9XkghTEbV9KlPS1xUsZ3u7/IQO4wxtcFB4bgpQPRcR3QCvezPcQ==", "license": "ISC" + }, + "node_modules/ws": { + "version": "8.21.3", + "resolved": "https://registry.npmjs.org/ws/-/ws-8.21.3.tgz", + "integrity": "sha512-201TZ/kPWxoPr/OKWjquZR1SWKXcvxdH+e1xrx89b3YbmzLMFCLfnaG1HFIgWzJOEWZ7MvpK++odZufgYR50Rw==", + "license": "MIT", + "engines": { + "node": ">=10.0.0" + }, + "peerDependencies": { + "bufferutil": "^4.0.1", + "utf-8-validate": ">=5.0.2" + }, + "peerDependenciesMeta": { + "bufferutil": { + "optional": true + }, + "utf-8-validate": { + "optional": true + } + } } } } diff --git a/lith/package.json b/lith/package.json index 6af9d6f9d3..8dd9af5c38 100644 --- a/lith/package.json +++ b/lith/package.json @@ -13,6 +13,7 @@ "dependencies": { "express": "^5.1.0", "mailparser": "^3.7.0", - "smtp-server": "^3.13.0" + "smtp-server": "^3.13.0", + "ws": "^8.21.3" } } diff --git a/lith/server.mjs b/lith/server.mjs index 9e012ef29a..337a6d1d57 100644 --- a/lith/server.mjs +++ b/lith/server.mjs @@ -47,6 +47,7 @@ if (typeof globalThis.awslambda === "undefined") { } import express from "express"; +import {attachMusicalSocket} from "./musical-socket.mjs"; import { sendStream } from "./stream-response.mjs"; import { userMediaTarget } from "./media-path.mjs"; import { readdirSync, readFileSync, existsSync, mkdirSync, writeFileSync, renameSync, statSync } from "fs"; @@ -1451,6 +1452,10 @@ if (DEV && HAS_SSL) { }); } +// Share authentication, concurrency guards and durable quotas with the HTTP lane. +const musicalAPI = await import(pathToFileURL(join(SYSTEM, "netlify/functions/easel-musical-jev.mjs")).href); +attachMusicalSocket(server, {authenticate:musicalAPI.authenticateMusical, decide:musicalAPI.musicalDecision}); + // --- Account deletions --- // Purges accounts whose grace period has ended (system/backend/ // account-deletion.mjs). A failed step backs off and resumes on a later run. diff --git a/system/backend/easel-musical-jev.mjs b/system/backend/easel-musical-jev.mjs index bbcaf24311..7e410378d2 100644 --- a/system/backend/easel-musical-jev.mjs +++ b/system/backend/easel-musical-jev.mjs @@ -16,21 +16,22 @@ export function musicalBudget(collection) { export function createMusicalHandler({authenticate,budget,evaluate=evaluateChoices,now=Date.now}={}) { const busy=new Set(); const reply=(statusCode,value)=>({statusCode,headers:{'Content-Type':'application/json','Cache-Control':'no-store','Access-Control-Allow-Origin':'*','Access-Control-Allow-Headers':'Content-Type, Authorization','Access-Control-Allow-Methods':'POST, OPTIONS'},body:JSON.stringify(value)}); - return async event=>{ + return async (event,{subject:trustedSubject,signal}={})=>{ if(event.httpMethod==='OPTIONS')return reply(200,{}); if(event.httpMethod!=='POST')return reply(405,{error:'POST only'}); if(!event.headers?.authorization)return reply(401,{error:'Sign in first'}); if(typeof event.body!=='string'||event.body.length>2048)return reply(400,{error:'Invalid observation'}); let body,request; try{body=JSON.parse(event.body);if(body.schema!=='walkieware-input/v1'||! /^[a-f0-9-]{36}$/i.test(body.sessionId)||!Number.isInteger(body.sequence)||body.sequence<1||body.sequence>1000||Object.keys(body).some(k=>!['schema','sessionId','sequence','features'].includes(k)))throw Error();request=musicalDecisionRequest(body.features);}catch{return reply(400,{error:'Invalid observation'});} - let subject;try{subject=await authenticate(event.headers);}catch{return reply(503,{error:'Account check unavailable'});} + let subject;try{subject=trustedSubject??await authenticate(event.headers);}catch{return reply(503,{error:'Account check unavailable'});} if(!subject)return reply(401,{error:'A valid account with a handle is required'}); if(busy.has(subject)||busy.size>=8)return reply(429,{error:'Decision already running'}); busy.add(subject); try{ if(!await budget.consume(subject,now()))return reply(429,{error:'Decision allowance reached'}); const started=performance.now(); - const result=await evaluate(request,{signal:AbortSignal.timeout(1000)}); + const deadline=AbortSignal.timeout(1000); + const result=await evaluate(request,{signal:signal?AbortSignal.any([signal,deadline]):deadline}); const answer=result.answers?.mapping, confidence=answer?.probabilities?.[answer.choice]; if(!Object.hasOwn(MUSICAL_CHOICES,answer?.choice)||!Number.isFinite(confidence)||confidence<0||confidence>1)throw Error('Invalid answer'); return reply(200,{schema:'walkieware-decision/v1',sessionId:body.sessionId,sequence:body.sequence,choice:answer.choice,confidence,elapsedMs:Math.round(performance.now()-started)}); diff --git a/system/netlify/functions/easel-musical-jev.mjs b/system/netlify/functions/easel-musical-jev.mjs index 43f7b64870..3d39de9cdf 100644 --- a/system/netlify/functions/easel-musical-jev.mjs +++ b/system/netlify/functions/easel-musical-jev.mjs @@ -3,8 +3,9 @@ import {authorize,getHandleOrEmail} from '../../backend/authorization.mjs'; import {createMusicalHandler,musicalBudget} from '../../backend/easel-musical-jev.mjs'; let storage; async function budget(){return storage??= (async()=>{const {db}=await connect();const c=db.collection('easel-musical-jev-budget');await c.createIndex({expiresAt:1},{expireAfterSeconds:0});return musicalBudget(c);})().catch(e=>{storage=null;throw e;});} -const handlerImpl=createMusicalHandler({ - authenticate:async headers=>{const user=await authorize(headers);if(!user?.sub)return null;const handle=await getHandleOrEmail(user.sub);return typeof handle==='string'&&handle.startsWith('@')?user.sub:null;}, +export const authenticateMusical=async headers=>{const user=await authorize(headers);if(!user?.sub)return null;const handle=await getHandleOrEmail(user.sub);return typeof handle==='string'&&handle.startsWith('@')?user.sub:null;}; +export const musicalDecision=createMusicalHandler({ + authenticate:authenticateMusical, budget:{consume:async(...args)=>(await budget()).consume(...args)} }); -export const handler=event=>handlerImpl(event); +export const handler=event=>musicalDecision(event);