diff --git a/apps/api/app/index.ts b/apps/api/app/index.ts index bb4bae4e..2979a97f 100644 --- a/apps/api/app/index.ts +++ b/apps/api/app/index.ts @@ -16,6 +16,7 @@ import { IndexApi } from './routes/index.js'; import { LibraryApi } from './routes/library.js'; import { MachineApi } from './routes/machine.js'; import { PairingCodeApi } from './routes/pairing-code.js'; +import { SessionApi } from './routes/session.js'; import { SteamApi } from './routes/steam.js'; import { UserApi } from './routes/user.js'; import { WaitlistApi } from './routes/waitlist.js'; @@ -44,6 +45,8 @@ const routes = app .route('/games', GameApi.route) .route('/pairing-code', PairingCodeApi.route) .route('/machine', MachineApi.route) + .route('/machine', SessionApi.machineRoute) + .route('/session', SessionApi.route) .route('/access-token', AccessTokenApi.route) .route('/waitlist', WaitlistApi.route) .onError((error, c) => { diff --git a/apps/api/app/routes/session.ts b/apps/api/app/routes/session.ts new file mode 100644 index 00000000..f84c9858 --- /dev/null +++ b/apps/api/app/routes/session.ts @@ -0,0 +1,337 @@ +import { Actor } from '@nestri/core/actor'; +import { Box } from '@nestri/core/box/index'; +import { ErrorCodes, VisibleError } from '@nestri/core/error'; +import { Examples } from '@nestri/core/examples'; +import { Game } from '@nestri/core/game/index'; +import { Identifier } from '@nestri/core/id'; +import { Session } from '@nestri/core/session/index'; +import { LinkedAccount } from '@nestri/core/user/linked-account'; +import { Hono } from 'hono'; +import { describeRoute } from 'hono-openapi'; +import { z } from 'zod'; + +import { ErrorResponses, machineOnly, notPublic, Result, validator } from '../utils'; + +/** + * Requesting a run, and carrying one out. + * + * Two very different callers meet on one resource here. A person asks for a + * run and then watches it; the host agent is handed the work and reports what + * happened. The rule that keeps them apart is that an agent may only see or + * touch a run whose box is placed on its own hardware, and it is enforced in + * the query rather than by the agent asking for its own work — a host + * credential is a long-lived secret on hardware in somebody's home, and what + * one leaking can reach is decided here. + */ +export namespace SessionApi { + /** + * One answer for "no such run" and "not your run". + * + * Both are the same refusal on purpose: an agent that could tell the + * difference could discover which ids exist by reporting states at them. + */ + function notYours(): never { + throw new VisibleError( + 'forbidden', + ErrorCodes.Permission.FORBIDDEN, + 'No such session, or it is not on this machine' + ); + } + + function conflict(message: string): never { + throw new VisibleError('already_exists', ErrorCodes.Validation.INVALID_STATE, message); + } + + /** The person a run belongs to, refusing a host acting as its owner. */ + function actingPerson(): string { + const actor = Actor.use(); + if (actor.type !== 'user' && actor.type !== 'member') { + throw new VisibleError( + 'forbidden', + ErrorCodes.Permission.INSUFFICIENT_PERMISSIONS, + 'Requesting or reading a session requires a user session' + ); + } + return actor.properties.userID; + } + + const StateReport = z + .object({ + state: Session.ReportableState.meta({ + description: 'Where the run has got to', + example: 'starting' + }), + errorMessage: z.string().max(1024).nullable().optional().meta({ + description: 'Why it failed. Kept only for a run that did', + example: Examples.Session.errorMessage + }) + }) + .strict(); + + export const route = new Hono() + .post( + '/', + notPublic, + describeRoute({ + tags: ['Session'], + summary: 'Ask for a run of a box', + description: + 'Creates the run in state `requested`, which is the work order the box’s host picks up. This makes no decision about where the run happens: a box already names the hardware it is placed on, so the run inherits it. Poll the run to watch it start, and re-read its ticket rather than keeping the first one.', + responses: { + 201: { + content: { 'application/json': { schema: Result(Session.Info) } }, + description: 'The run has been requested' + }, + 400: ErrorResponses[400], + 401: ErrorResponses[401], + 403: ErrorResponses[403], + 404: ErrorResponses[404], + 409: ErrorResponses[409] + } + }), + validator( + 'json', + z + .object({ + boxId: z.string().min(1).meta({ + description: 'The box to run', + example: Examples.Session.boxId + }), + gameId: z.string().min(1).meta({ + description: 'The game to launch', + example: Examples.Session.gameId + }), + linkedAccountId: z.string().min(1).optional().meta({ + description: + 'Which linked account is playing. Defaults to the one the caller signed in with', + example: Examples.Session.linkedAccountId + }) + }) + // Strict, so that naming hardware is a validation error rather + // than a field quietly ignored. There is nothing to choose: + // asking for a run is not where a box is placed. + .strict() + ), + async (c) => { + const body = c.req.valid('json'); + const userId = actingPerson(); + + const box = await Box.fromID(body.boxId); + if (!box || box.userId !== userId) { + // Somebody else's box and a box that was never created are the + // same answer, so ids cannot be probed for. + throw new VisibleError( + 'not_found', + ErrorCodes.NotFound.RESOURCE_NOT_FOUND, + 'No such box, or it is not yours' + ); + } + + const game = await Game.fromID(body.gameId); + if (!game) { + throw new VisibleError( + 'not_found', + ErrorCodes.NotFound.RESOURCE_NOT_FOUND, + 'No such game' + ); + } + + const actor = Actor.use(); + const linkedAccountId = + body.linkedAccountId || + (actor.type === 'user' ? actor.properties.linkedAccountID : '') || + ''; + if (!linkedAccountId) { + // Which account is playing is the question the "who's playing?" + // screen asks, and some credentials carry no answer to it. Then + // the caller has to say. + throw new VisibleError( + 'validation', + ErrorCodes.Validation.MISSING_REQUIRED_FIELD, + 'Say which linked account is playing', + 'linkedAccountId' + ); + } + const linked = await LinkedAccount.fromID(linkedAccountId); + if (!linked || linked.userId !== userId) { + throw new VisibleError( + 'forbidden', + ErrorCodes.Permission.FORBIDDEN, + 'That account is not linked to you' + ); + } + + // A box runs one thing at a time. Refusing is the honest answer; + // starting a second run would leave two rows that both think they + // own the same hardware. + const active = await Session.activeForBox(box.id); + if (active) { + conflict('That box already has a run that has not stopped'); + } + + const session = await Session.create({ + id: Identifier.ascending('session'), + boxId: box.id, + gameId: game.id, + linkedAccountId + }); + return c.json({ data: session }, 201); + } + ) + .get( + '/:id', + notPublic, + describeRoute({ + tags: ['Session'], + summary: 'Read a run you asked for', + description: + 'Poll this while a run starts. The ticket appears part-way through and is republished as addresses are discovered, so re-read it rather than keeping the first one — a client that treats the first ticket as final works on a local network and fails from anywhere else.', + responses: { + 200: { + content: { 'application/json': { schema: Result(Session.Info) } }, + description: 'The run as it stands' + }, + 401: ErrorResponses[401], + 403: ErrorResponses[403], + 404: ErrorResponses[404] + } + }), + validator( + 'param', + z.object({ + id: z.string().meta({ description: 'The run to read', example: Examples.Session.id }) + }) + ), + async (c) => { + const session = await Session.forOwner({ + id: c.req.valid('param').id, + userId: actingPerson() + }); + if (!session) { + // Owner-scoped in the query, so somebody else's run and one that + // never existed answer the same way. + throw new VisibleError( + 'not_found', + ErrorCodes.NotFound.RESOURCE_NOT_FOUND, + 'No such session, or it is not yours' + ); + } + return c.json({ data: session }); + } + ) + .post( + '/:id/state', + machineOnly, + describeRoute({ + tags: ['Session'], + summary: 'Report where a run has got to', + description: + 'For the host the run’s box is placed on, and no other. Moving a run out of `requested` is the claim, and it is a compare-and-set: exactly one caller can take a given run, and one that loses gets 409. Re-reporting a state already reported is fine and changes nothing, including the timestamps a run is billed on. A transition that does not exist is 409 and the run does not move.', + responses: { + 200: { + content: { 'application/json': { schema: Result(Session.Info) } }, + description: 'The run as it stands after the report' + }, + 400: ErrorResponses[400], + 403: ErrorResponses[403], + 409: ErrorResponses[409] + } + }), + validator('param', z.object({ id: z.string() })), + validator('json', StateReport), + async (c) => { + const body = c.req.valid('json'); + const result = await Session.transition({ + id: c.req.valid('param').id, + machineId: Actor.machineID, + state: body.state, + errorMessage: body.errorMessage ?? null + }); + + switch (result.outcome) { + case 'forbidden': + notYours(); + case 'illegal': + conflict(`A run in state ${result.session?.state} cannot become ${body.state}`); + case 'lost': + conflict('Another caller moved this run first'); + default: + // `moved` and `unchanged` are both success. An agent retrying + // after a lost response must not be told it broke something. + return c.json({ data: result.session }); + } + } + ) + .post( + '/:id/ticket', + machineOnly, + describeRoute({ + tags: ['Session'], + summary: 'Publish the address a client should connect to', + description: + 'For the host the run’s box is placed on, and no other. Republish freely: a later ticket is a better address for the same run, not a second run, and the address changes as more of them are discovered. A run that has stopped has no address, so that is 409.', + responses: { + 200: { + content: { 'application/json': { schema: Result(Session.Info) } }, + description: 'The ticket is published' + }, + 400: ErrorResponses[400], + 403: ErrorResponses[403], + 409: ErrorResponses[409] + } + }), + validator('param', z.object({ id: z.string() })), + validator( + 'json', + z + .object({ + ticket: z.string().min(1).meta({ + description: 'The current connect ticket', + example: Examples.Session.ticket + }) + }) + .strict() + ), + async (c) => { + const result = await Session.publishTicket({ + id: c.req.valid('param').id, + machineId: Actor.machineID, + ticket: c.req.valid('json').ticket + }); + + switch (result.outcome) { + case 'forbidden': + notYours(); + case 'closed': + conflict('That run has stopped, so it has no address to publish'); + default: + return c.json({ data: result.session }); + } + } + ); + + /** + * The host agent's side of the same resource, mounted where a host looks + * for it: everything a box asks about itself lives under one prefix. + */ + export const machineRoute = new Hono().get( + '/jobs', + machineOnly, + describeRoute({ + tags: ['Session'], + summary: 'Ask for work', + description: + 'Returns the runs waiting to be started on the calling host, and only those — the host comes from its own credentials and the scope is the query, so a box cannot see work for another. Poll at the cadence the heartbeat hands down. Each job carries its kind, so a second kind of work is an addition rather than a change of shape.', + responses: { + 200: { + content: { 'application/json': { schema: Result(z.array(Session.Job)) } }, + description: 'Work waiting for this host, oldest first' + }, + 403: ErrorResponses[403] + } + }), + async (c) => { + return c.json({ data: await Session.listJobsForMachine(Actor.machineID) }); + } + ); +} diff --git a/apps/api/test/session.test.ts b/apps/api/test/session.test.ts new file mode 100644 index 00000000..03c6ba5f --- /dev/null +++ b/apps/api/test/session.test.ts @@ -0,0 +1,564 @@ +import { afterAll, describe, expect, test } from 'bun:test'; + +import { AccessToken } from '@nestri/core/access-token/index'; +import { Box } from '@nestri/core/box/index'; +import { Fixtures } from '@nestri/core/db/fixtures'; +import { testDb } from '@nestri/core/db/test'; +import { Game } from '@nestri/core/game/index'; +import { Identifier } from '@nestri/core/id'; +import { Machine } from '@nestri/core/machine/index'; +import { Session } from '@nestri/core/session/index'; + +import { app } from '../app/index'; +import './setup'; + +const sql = testDb(); + +const createdUserIds: string[] = []; +const createdGameIds: string[] = []; + +async function newGame(steamAppId: number): Promise { + const [row] = await Game.upsert({ + id: Identifier.ascending('game'), + steamAppId, + slug: `session-route-${steamAppId}`, + name: `Session Route ${steamAppId}` + }); + if (!row) throw new Error('expected a game row'); + createdGameIds.push(row.id); + return row.id; +} + +/** + * Everything one session needs, plus both sets of credentials that reach it. + * + * The person authenticates with a personal token, which is the one user + * credential a test can mint without an auth service; the host authenticates + * as itself with the secret registration hands back exactly once. + */ +async function scene(label: string, steamAppId: number) { + const owner = await Fixtures.owner(label); + createdUserIds.push(owner.userId); + + const registered = await Machine.register({ + id: Identifier.ascending('machine'), + ownerUserId: owner.userId, + teamId: owner.teamId, + label + }); + + const box = await Box.create({ + id: Identifier.ascending('box'), + userId: owner.userId, + machineId: registered.id, + label, + tier: 'sm' + }); + + const pat = await AccessToken.create({ + id: Identifier.ascending('accessToken'), + ownerUserId: owner.userId, + // Null on purpose: a token scoped to the user alone makes the caller a + // plain user actor, which is the credential a person browsing has. + teamId: null, + name: label + }); + + return { + owner, + box, + machineId: registered.id, + gameId: await newGame(steamAppId), + user: { + authorization: `Bearer ${pat.token}`, + 'content-type': 'application/json' + } as Record, + host: { + 'x-nestri-machine-id': registered.id, + 'x-nestri-machine-secret': registered.secret, + 'content-type': 'application/json' + } as Record + }; +} + +async function requestSession(s: Awaited>) { + const res = await app.request('/session', { + method: 'POST', + headers: s.user, + body: JSON.stringify({ + boxId: s.box.id, + gameId: s.gameId, + linkedAccountId: s.owner.linkedAccountId + }) + }); + const body = (await res.json()) as any; + return { res, body }; +} + +afterAll(async () => { + if (createdUserIds.length > 0) { + await sql`delete from "box" where user_id in ${sql(createdUserIds)}`; + await sql`delete from "user" where id in ${sql(createdUserIds)}`; + createdUserIds.length = 0; + } + if (createdGameIds.length > 0) { + await sql`delete from "game" where id in ${sql(createdGameIds)}`; + createdGameIds.length = 0; + } +}); + +describe('POST /session', () => { + test('a request creates the job, in the envelope both ends read', async () => { + const s = await scene('route-create', 5500); + const { res, body } = await requestSession(s); + + expect(res.status).toBe(201); + // The field names are the contract. A rename on either side produces a + // host that starts, reads nothing, and reports success — so the shape + // is asserted whole rather than field by field. + expect(Object.keys(body)).toEqual(['data']); + expect(body.data).toEqual({ + id: body.data.id, + boxId: s.box.id, + gameId: s.gameId, + linkedAccountId: s.owner.linkedAccountId, + state: 'requested', + ticket: null, + timeStarted: null, + timeStopped: null, + errorMessage: null + }); + expect(body.data.id.startsWith('ses_')).toBe(true); + }); + + test('creating a session makes no placement decision', async () => { + const s = await scene('route-noplacement', 5501); + const { body } = await requestSession(s); + + // A session inherits its machine through its box, so there is nothing + // to choose here and no way for a caller to ask for a host. + expect(body.data).not.toHaveProperty('machineId'); + + const withHost = await app.request('/session', { + method: 'POST', + headers: s.user, + body: JSON.stringify({ + boxId: s.box.id, + gameId: s.gameId, + linkedAccountId: s.owner.linkedAccountId, + machineId: s.machineId + }) + }); + expect(withHost.status).toBe(400); + }); + + test('a box somebody else owns is not there to run', async () => { + const mine = await scene('route-mine', 5502); + const theirs = await scene('route-theirs', 5503); + + const res = await app.request('/session', { + method: 'POST', + headers: mine.user, + body: JSON.stringify({ + boxId: theirs.box.id, + gameId: mine.gameId, + linkedAccountId: mine.owner.linkedAccountId + }) + }); + expect(res.status).toBe(404); + + const unknown = await app.request('/session', { + method: 'POST', + headers: mine.user, + body: JSON.stringify({ + boxId: Identifier.ascending('box'), + gameId: mine.gameId, + linkedAccountId: mine.owner.linkedAccountId + }) + }); + // Owner-scoped in the query, so somebody else's box and a box that was + // never created are the same answer. + expect(unknown.status).toBe(404); + expect(await res.json()).toEqual(await unknown.json()); + }); + + test('a box already running refuses a second run rather than picking one', async () => { + const s = await scene('route-busy', 5504); + expect((await requestSession(s)).res.status).toBe(201); + + const second = await requestSession(s); + expect(second.res.status).toBe(409); + expect(second.body.type).toBe('already_exists'); + }); + + test('you can only play as an account you have linked', async () => { + const mine = await scene('route-account-mine', 5505); + const theirs = await scene('route-account-theirs', 5506); + + const res = await app.request('/session', { + method: 'POST', + headers: mine.user, + body: JSON.stringify({ + boxId: mine.box.id, + gameId: mine.gameId, + linkedAccountId: theirs.owner.linkedAccountId + }) + }); + expect(res.status).toBe(403); + }); + + test('an unknown game is a 404 and not a foreign key crash', async () => { + const s = await scene('route-nogame', 5507); + const res = await app.request('/session', { + method: 'POST', + headers: s.user, + body: JSON.stringify({ + boxId: s.box.id, + gameId: Identifier.ascending('game'), + linkedAccountId: s.owner.linkedAccountId + }) + }); + expect(res.status).toBe(404); + }); + + test('a host cannot ask for a session on its owner’s behalf', async () => { + const s = await scene('route-hostcreate', 5508); + const res = await app.request('/session', { + method: 'POST', + headers: s.host, + body: JSON.stringify({ + boxId: s.box.id, + gameId: s.gameId, + linkedAccountId: s.owner.linkedAccountId + }) + }); + // A box holds credentials but is not the person who owns it. + expect(res.status).toBe(403); + }); + + test('requesting a session requires a signed-in person', async () => { + const res = await app.request('/session', { + method: 'POST', + headers: { 'content-type': 'application/json' }, + body: JSON.stringify({ boxId: 'box_x', gameId: 'gam_x', linkedAccountId: 'lac_x' }) + }); + expect(res.status).toBe(401); + }); +}); + +describe('GET /session/:id', () => { + test('the owner reads their own run, ticket and all', async () => { + const s = await scene('route-read', 5510); + const { body } = await requestSession(s); + + await app.request(`/session/${body.data.id}/state`, { + method: 'POST', + headers: s.host, + body: JSON.stringify({ state: 'starting' }) + }); + await app.request(`/session/${body.data.id}/ticket`, { + method: 'POST', + headers: s.host, + body: JSON.stringify({ ticket: 'nodeaaa-one' }) + }); + + const res = await app.request(`/session/${body.data.id}`, { headers: s.user }); + expect(res.status).toBe(200); + const read = (await res.json()) as any; + expect(read.data.state).toBe('starting'); + // A ticket may appear while the state is still `starting`, and the + // client is expected to re-read rather than cache the first one. + expect(read.data.ticket).toBe('nodeaaa-one'); + }); + + test('somebody else’s run is not visible, and neither is its absence', async () => { + const mine = await scene('route-read-mine', 5511); + const theirs = await scene('route-read-theirs', 5512); + const { body } = await requestSession(theirs); + + const forbidden = await app.request(`/session/${body.data.id}`, { headers: mine.user }); + const unknown = await app.request(`/session/${Identifier.ascending('session')}`, { + headers: mine.user + }); + expect(forbidden.status).toBe(404); + expect(unknown.status).toBe(404); + expect(await forbidden.json()).toEqual(await unknown.json()); + }); + + test('reading a run requires a signed-in person', async () => { + const res = await app.request('/session/ses_whatever'); + expect(res.status).toBe(401); + }); +}); + +describe('GET /machine/jobs', () => { + test('a host is handed the work for its own boxes, with the kind on the wire', async () => { + const s = await scene('route-jobs', 5520); + const { body } = await requestSession(s); + + const res = await app.request('/machine/jobs', { headers: s.host }); + expect(res.status).toBe(200); + const jobs = (await res.json()) as any; + expect(Object.keys(jobs)).toEqual(['data']); + expect(jobs.data).toHaveLength(1); + expect(jobs.data[0]).toEqual({ + kind: 'session.start', + sessionId: body.data.id, + boxId: s.box.id, + boxTier: 'sm', + gameId: s.gameId, + steamAppId: 5520, + linkedAccountId: s.owner.linkedAccountId + }); + }); + + test('a host never sees work for a box on other hardware', async () => { + const mine = await scene('route-jobs-mine', 5521); + const theirs = await scene('route-jobs-theirs', 5522); + await requestSession(theirs); + + const res = await app.request('/machine/jobs', { headers: mine.host }); + expect(res.status).toBe(200); + // Scoped in the query rather than by the host asking for its own work. + expect(((await res.json()) as any).data).toEqual([]); + }); + + test('bad credentials are indistinguishable from none', async () => { + const s = await scene('route-jobs-auth', 5523); + const wrong = await app.request('/machine/jobs', { + headers: { ...s.host, 'x-nestri-machine-secret': 'msk_wrong' } + }); + const none = await app.request('/machine/jobs'); + expect(wrong.status).toBe(403); + expect(none.status).toBe(403); + // Bad credentials fall through to public and are then forbidden, so + // probing tells an attacker nothing. Asserting the two are identical is + // the only way that stays true. + expect(await wrong.json()).toEqual(await none.json()); + }); + + test('a person cannot poll for jobs', async () => { + const s = await scene('route-jobs-person', 5524); + const res = await app.request('/machine/jobs', { headers: s.user }); + expect(res.status).toBe(403); + }); +}); + +describe('POST /session/:id/state', () => { + test('the claim moves the row, and the job stops being offered', async () => { + const s = await scene('route-claim', 5530); + const { body } = await requestSession(s); + + const res = await app.request(`/session/${body.data.id}/state`, { + method: 'POST', + headers: s.host, + body: JSON.stringify({ state: 'starting' }) + }); + expect(res.status).toBe(200); + expect(((await res.json()) as any).data.state).toBe('starting'); + + const jobs = await app.request('/machine/jobs', { headers: s.host }); + expect(((await jobs.json()) as any).data).toEqual([]); + }); + + test('the same host re-reporting a state it already reported is fine', async () => { + const s = await scene('route-claim-retry', 5531); + const { body } = await requestSession(s); + + const report = () => + app.request(`/session/${body.data.id}/state`, { + method: 'POST', + headers: s.host, + body: JSON.stringify({ state: 'starting' }) + }); + + expect((await report()).status).toBe(200); + // An agent retrying after a lost response must not be told it broke + // something. + const again = await report(); + expect(again.status).toBe(200); + expect(((await again.json()) as any).data.state).toBe('starting'); + }); + + test('a different host reporting anything is refused, and learns nothing', async () => { + const mine = await scene('route-claim-mine', 5532); + const theirs = await scene('route-claim-theirs', 5533); + const { body } = await requestSession(theirs); + + const other = await app.request(`/session/${body.data.id}/state`, { + method: 'POST', + headers: mine.host, + body: JSON.stringify({ state: 'starting' }) + }); + const unknown = await app.request(`/session/${Identifier.ascending('session')}/state`, { + method: 'POST', + headers: mine.host, + body: JSON.stringify({ state: 'starting' }) + }); + + expect(other.status).toBe(403); + expect(unknown.status).toBe(403); + expect(await other.json()).toEqual(await unknown.json()); + expect((await Session.fromID(body.data.id))?.state).toBe('requested'); + }); + + test('a transition that is not allowed is a conflict, and the row stays put', async () => { + const s = await scene('route-claim-illegal', 5534); + const { body } = await requestSession(s); + + const skipped = await app.request(`/session/${body.data.id}/state`, { + method: 'POST', + headers: s.host, + body: JSON.stringify({ state: 'live' }) + }); + expect(skipped.status).toBe(409); + expect((await Session.fromID(body.data.id))?.state).toBe('requested'); + }); + + test('a stopped run cannot be started again', async () => { + const s = await scene('route-claim-terminal', 5535); + const { body } = await requestSession(s); + const report = (state: string, errorMessage?: string) => + app.request(`/session/${body.data.id}/state`, { + method: 'POST', + headers: s.host, + body: JSON.stringify({ state, errorMessage }) + }); + + expect((await report('starting')).status).toBe(200); + expect((await report('failed', 'the guest never came up')).status).toBe(200); + expect((await report('starting')).status).toBe(409); + + const failed = await Session.fromID(body.data.id); + expect(failed?.state).toBe('failed'); + expect(failed?.errorMessage).toBe('the guest never came up'); + }); + + test('a duplicate live report does not extend a run somebody is billed for', async () => { + const s = await scene('route-claim-billing', 5536); + const { body } = await requestSession(s); + const report = (state: string) => + app.request(`/session/${body.data.id}/state`, { + method: 'POST', + headers: s.host, + body: JSON.stringify({ state }) + }); + + await report('starting'); + const live = (await (await report('live')).json()) as any; + expect(live.data.timeStarted).not.toBeNull(); + + const again = (await (await report('live')).json()) as any; + expect(again.data.timeStarted).toBe(live.data.timeStarted); + }); + + test('a state nobody defined is a validation error, not a conflict', async () => { + const s = await scene('route-claim-bogus', 5537); + const { body } = await requestSession(s); + const res = await app.request(`/session/${body.data.id}/state`, { + method: 'POST', + headers: s.host, + body: JSON.stringify({ state: 'exploded' }) + }); + expect(res.status).toBe(400); + }); + + test('a person cannot report a state on their own session', async () => { + const s = await scene('route-claim-person', 5538); + const { body } = await requestSession(s); + const res = await app.request(`/session/${body.data.id}/state`, { + method: 'POST', + headers: s.user, + body: JSON.stringify({ state: 'starting' }) + }); + // Terminal states are written by the agent alone; a person closing the + // app is not the same fact as a run that stopped. + expect(res.status).toBe(403); + }); +}); + +describe('POST /session/:id/ticket', () => { + test('a later ticket replaces the first, because it is a better address', async () => { + const s = await scene('route-ticket', 5540); + const { body } = await requestSession(s); + await app.request(`/session/${body.data.id}/state`, { + method: 'POST', + headers: s.host, + body: JSON.stringify({ state: 'starting' }) + }); + + const publish = (ticket: string) => + app.request(`/session/${body.data.id}/ticket`, { + method: 'POST', + headers: s.host, + body: JSON.stringify({ ticket }) + }); + + const first = await publish('nodeaaa-one'); + expect(first.status).toBe(200); + expect(((await first.json()) as any).data.ticket).toBe('nodeaaa-one'); + + const second = await publish('nodeaaa-two'); + expect(((await second.json()) as any).data.ticket).toBe('nodeaaa-two'); + expect(await Session.listByBox(s.box.id)).toHaveLength(1); + }); + + test('a different host cannot publish an address for someone else’s run', async () => { + const mine = await scene('route-ticket-mine', 5541); + const theirs = await scene('route-ticket-theirs', 5542); + const { body } = await requestSession(theirs); + + const res = await app.request(`/session/${body.data.id}/ticket`, { + method: 'POST', + headers: mine.host, + body: JSON.stringify({ ticket: 'nodeaaa-stolen' }) + }); + expect(res.status).toBe(403); + expect((await Session.fromID(body.data.id))?.ticket).toBeNull(); + }); + + test('a stopped run has no address to publish', async () => { + const s = await scene('route-ticket-dead', 5543); + const { body } = await requestSession(s); + const report = (state: string) => + app.request(`/session/${body.data.id}/state`, { + method: 'POST', + headers: s.host, + body: JSON.stringify({ state }) + }); + await report('starting'); + await report('live'); + await report('ended'); + + const res = await app.request(`/session/${body.data.id}/ticket`, { + method: 'POST', + headers: s.host, + body: JSON.stringify({ ticket: 'nodeaaa-late' }) + }); + expect(res.status).toBe(409); + expect((await Session.fromID(body.data.id))?.ticket).toBeNull(); + }); + + test('a ticket has to say something', async () => { + const s = await scene('route-ticket-empty', 5544); + const { body } = await requestSession(s); + const res = await app.request(`/session/${body.data.id}/ticket`, { + method: 'POST', + headers: s.host, + body: JSON.stringify({ ticket: '' }) + }); + expect(res.status).toBe(400); + }); +}); + +describe('Session routes in the spec', () => { + test('every path a caller needs is documented', async () => { + const res = await app.request('/doc'); + const paths = Object.keys(((await res.json()) as any).paths); + expect(paths).toContain('/session'); + expect(paths).toContain('/session/{id}'); + expect(paths).toContain('/session/{id}/state'); + expect(paths).toContain('/session/{id}/ticket'); + expect(paths).toContain('/machine/jobs'); + }); +}); diff --git a/packages/core/src/box/index.ts b/packages/core/src/box/index.ts index fd7926fc..bfcd36c7 100644 --- a/packages/core/src/box/index.ts +++ b/packages/core/src/box/index.ts @@ -5,6 +5,7 @@ import { Database } from '../db/index.js'; import { Examples } from '../examples.js'; import { fn } from '../fn.js'; import { BoxState, BoxTable, BoxTier } from './box.sql.js'; +import { Placement } from './placement.js'; /** * A VM someone owns. @@ -77,6 +78,29 @@ export namespace Box { } ); + /** + * Create a box and let something else decide where it runs. + * + * The placement seam is here, at creation, and nowhere else: `machineId` is + * set once and every later question about which hardware a box — or a run + * of it — belongs to is answered by joining through this row. A caller that + * knows the host still uses `create`; a caller acting for a person does not + * know and must not guess, which is what this overload is for. + */ + export const createPlaced = async ( + input: { id: string; userId: string; label: string; tier: Info['tier'] }, + placer?: Placement.Placer + ) => { + const machineId = await Placement.choose({ userId: input.userId, tier: input.tier }, placer); + return create({ + id: input.id, + userId: input.userId, + machineId, + label: input.label, + tier: input.tier + }); + }; + export const fromID = fn(Info.shape.id, async (id) => { return Database.use(async (tx) => { return tx diff --git a/packages/core/src/box/placement.test.ts b/packages/core/src/box/placement.test.ts new file mode 100644 index 00000000..77867198 --- /dev/null +++ b/packages/core/src/box/placement.test.ts @@ -0,0 +1,90 @@ +import { afterAll, describe, expect, test } from 'bun:test'; + +import { Fixtures } from '../db/fixtures.js'; +import { testDb } from '../db/test.js'; +import { Identifier } from '../id.js'; +import { Placement } from './placement.js'; +import { Box } from './index.js'; + +const sql = testDb(); + +const createdUserIds: string[] = []; + +async function newOwner(label: string) { + const o = await Fixtures.owner(label); + createdUserIds.push(o.userId); + return o; +} + +afterAll(async () => { + if (createdUserIds.length > 0) { + await sql`delete from "box" where user_id in ${sql(createdUserIds)}`; + await sql`delete from "user" where id in ${sql(createdUserIds)}`; + createdUserIds.length = 0; + } +}); + +describe('Placement', () => { + test('a box is placed when it is created, and the caller names no host', async () => { + const owner = await newOwner('place-one'); + const machineId = await Fixtures.machine(owner, 'place-one-host'); + + // No `machineId` in the input: choosing the host is the placer's job, + // and the whole point of the interface is that the caller cannot do it. + const box = await Box.createPlaced({ + id: Identifier.ascending('box'), + userId: owner.userId, + label: 'living room', + tier: 'sm' + }); + + expect(box.machineId).toBe(machineId); + }); + + test('nowhere to put it is an answer, not a crash', async () => { + const owner = await newOwner('place-none'); + + // `box.machineId` is notNull, so a placer with no candidate must refuse + // rather than hand back something the insert would reject. + await expect( + Placement.choose({ userId: owner.userId, tier: 'sm' }) + ).rejects.toThrow(); + }); + + test('more than one candidate is refused rather than picked silently', async () => { + const owner = await newOwner('place-two'); + await Fixtures.machine(owner, 'place-two-a'); + await Fixtures.machine(owner, 'place-two-b'); + + // There is no policy for choosing between hosts yet. Inventing one here + // is how a placement decision ends up buried in the caller: the refusal + // is what keeps the choice in one replaceable place. + await expect( + Placement.choose({ userId: owner.userId, tier: 'sm' }) + ).rejects.toThrow(); + }); + + test('the placer is swappable without touching box creation', async () => { + const owner = await newOwner('place-swap'); + const machineId = await Fixtures.machine(owner, 'place-swap-host'); + + const asked: unknown[] = []; + const box = await Box.createPlaced( + { + id: Identifier.ascending('box'), + userId: owner.userId, + label: 'bedroom', + tier: 'lg' + }, + async (input) => { + asked.push(input); + return machineId; + } + ); + + expect(box.machineId).toBe(machineId); + // The placer is told who the box is for and what size was asked for, + // which is the whole input a real scheduler needs. + expect(asked).toEqual([{ userId: owner.userId, tier: 'lg' }]); + }); +}); diff --git a/packages/core/src/box/placement.ts b/packages/core/src/box/placement.ts new file mode 100644 index 00000000..64f7d972 --- /dev/null +++ b/packages/core/src/box/placement.ts @@ -0,0 +1,73 @@ +import z from 'zod'; + +import { ErrorCodes, VisibleError } from '../error.js'; +import { Machine } from '../machine/index.js'; +import { BoxTier } from './box.sql.js'; + +/** + * Deciding which host a box runs on. + * + * This is a seam and not an algorithm. `box.machineId` is set once, when the + * box is created, and everything downstream — a session, its job, the state + * reports that follow — reaches the right hardware by joining through the box. + * So there is exactly one moment where placement happens, and the value of + * naming it now is that a real scheduler replaces this file and nothing else. + * + * The wrong shape, and the tempting one, is to place a box when a *run* is + * requested. That spreads the decision across every caller that starts + * something and leaves nowhere to put a scheduler later. + */ +export namespace Placement { + export const Request = z.object({ + userId: z.string().meta({ description: 'Who the box is for' }), + tier: z.enum(BoxTier.enumValues).meta({ description: 'The size that was asked for' }) + }); + + export type Request = z.infer; + + /** + * Answers "which host should run this box?" with a machine id. + * + * Asynchronous and allowed to refuse: capacity is a real answer, and a + * placer that cannot honour a request must say so rather than return + * something the insert would reject — `box.machineId` is not nullable. + */ + export type Placer = (request: Request) => Promise; + + /** + * The implementation there is hardware for: place it on the caller's host. + * + * Deliberately refuses when the answer is not forced. With no host there is + * nothing to place on; with several there is a choice to make and no policy + * to make it with, and picking the first row would be a scheduling decision + * taken by accident and impossible to find later. Refusing keeps the choice + * in this one function. + */ + export const onlyHost: Placer = async (request) => { + const hosts = await Machine.listByOwner(request.userId); + + if (hosts.length === 0) { + throw new VisibleError( + 'not_found', + ErrorCodes.NotFound.RESOURCE_NOT_FOUND, + 'You have no registered host to run a box on' + ); + } + if (hosts.length > 1) { + // Not a caller error: the request is fine and the system cannot yet + // answer it. An orchestrator is what closes this. todo(d-0048) + throw new VisibleError( + 'internal', + ErrorCodes.Server.SERVICE_UNAVAILABLE, + 'More than one host could run this box, and choosing between them is not supported yet' + ); + } + + return hosts[0]!.id; + }; + + /** Place a box, using `onlyHost` unless a caller supplies its own placer. */ + export async function choose(request: Request, placer: Placer = onlyHost): Promise { + return placer(Request.parse(request)); + } +} diff --git a/packages/core/src/session/index.ts b/packages/core/src/session/index.ts index fbb0f3f9..aa3c9061 100644 --- a/packages/core/src/session/index.ts +++ b/packages/core/src/session/index.ts @@ -1,9 +1,11 @@ -import { and, desc, eq, isNull, sql } from 'drizzle-orm'; +import { and, desc, eq, inArray, isNull, notInArray, sql } from 'drizzle-orm'; import z from 'zod'; +import { BoxTable, BoxTier } from '../box/box.sql.js'; import { Database } from '../db/index.js'; import { Examples } from '../examples.js'; import { fn } from '../fn.js'; +import { GameTable } from '../game/game.sql.js'; import { SessionState, SessionTable } from './session.sql.js'; /** @@ -188,6 +190,328 @@ export namespace Session { } ); + /** + * The states an agent is allowed to move a run into. + * + * `requested` is missing on purpose: it is written once, when the row is + * created, and nothing may put a run back there. + */ + export const ReportableState = z.enum(['starting', 'live', 'ended', 'failed']); + + export type ReportableState = z.infer; + + /** + * Where a run may go next, and nowhere else. + * + * `requested → live` is missing although it is the tempting shortcut: + * skipping `starting` means nothing ever holds the claim, and the claim is + * the only mutual exclusion in this design. `ended` and `failed` are + * terminal, so their entries are empty rather than absent — a state with no + * exits is a fact worth writing down. + */ + export const NEXT_STATES: Record = { + requested: ['starting'], + starting: ['live', 'failed'], + live: ['ended', 'failed'], + ended: [], + failed: [] + }; + + /** + * A unit of work handed to the agent that will carry it out. + * + * There is no queue: a run in state `requested` *is* the work order, and + * the agent that fulfils it moves that same row along. Two sources of truth + * for one piece of work is how a queue and a database come to disagree + * about whether something ran. + * + * `kind` is on the wire while there is only one value, so that a second + * kind is an addition rather than a redesign of the poll. + */ + export const Job = z + .object({ + kind: z.literal('session.start').meta({ + description: 'What the agent is being asked to do', + example: 'session.start' + }), + sessionId: z.string().meta({ + description: 'The run to report progress against', + example: Examples.Session.id + }), + boxId: z.string().meta({ + description: 'The box to start', + example: Examples.Session.boxId + }), + boxTier: z.enum(BoxTier.enumValues).meta({ + description: 'The size the box was asked for, which also sets output geometry', + example: Examples.Box.tier + }), + gameId: z.string().meta({ + description: 'The game to launch', + example: Examples.Session.gameId + }), + steamAppId: z.number().int().meta({ + description: 'The same game, in the id the store knows it by', + example: Examples.Game.steamAppId + }), + linkedAccountId: z.string().meta({ + description: 'Which linked account is playing', + example: Examples.Session.linkedAccountId + }) + }) + .meta({ + ref: 'SessionJob', + description: 'One run waiting to be started, as handed to the agent that will start it' + }); + + export type Job = z.infer; + + /** + * The work waiting for one host. + * + * The scope is the join and not a filter the caller asks for: a box names + * the hardware it is placed on, a run reaches its hardware through its box, + * and so what one set of long-lived credentials can see is decided by this + * `where` clause rather than by whoever is holding them. + */ + export const listJobsForMachine = fn(z.string(), async (machineId) => { + return Database.use(async (tx) => { + return tx + .select({ session: SessionTable, box: BoxTable, game: GameTable }) + .from(SessionTable) + .innerJoin(BoxTable, eq(SessionTable.boxId, BoxTable.id)) + .innerJoin(GameTable, eq(SessionTable.gameId, GameTable.id)) + .where( + and( + eq(BoxTable.machineId, machineId), + eq(SessionTable.state, 'requested'), + isNull(SessionTable.timeDeleted), + isNull(BoxTable.timeDeleted) + ) + ) + .orderBy(SessionTable.timeCreated) + .then((rows) => + rows.map( + (row): Job => ({ + kind: 'session.start', + sessionId: row.session.id, + boxId: row.box.id, + boxTier: row.box.tier as Job['boxTier'], + gameId: row.game.id, + steamAppId: row.game.steamAppId, + linkedAccountId: row.session.linkedAccountId + }) + ) + ); + }); + }); + + /** One run, visible only to the host its box is placed on. */ + export const forMachine = fn( + z.object({ id: Info.shape.id, machineId: z.string() }), + async (input) => { + return Database.use(async (tx) => { + return tx + .select({ session: SessionTable }) + .from(SessionTable) + .innerJoin(BoxTable, eq(SessionTable.boxId, BoxTable.id)) + .where( + and( + eq(SessionTable.id, input.id), + eq(BoxTable.machineId, input.machineId), + isNull(SessionTable.timeDeleted), + isNull(BoxTable.timeDeleted) + ) + ) + .then((rows) => { + const row = rows.at(0); + return row ? serialize(row.session) : null; + }); + }); + } + ); + + /** One run, visible only to the person who owns its box. */ + export const forOwner = fn(z.object({ id: Info.shape.id, userId: z.string() }), async (input) => { + return Database.use(async (tx) => { + return tx + .select({ session: SessionTable }) + .from(SessionTable) + .innerJoin(BoxTable, eq(SessionTable.boxId, BoxTable.id)) + .where( + and( + eq(SessionTable.id, input.id), + eq(BoxTable.userId, input.userId), + isNull(SessionTable.timeDeleted), + isNull(BoxTable.timeDeleted) + ) + ) + .then((rows) => { + const row = rows.at(0); + return row ? serialize(row.session) : null; + }); + }); + }); + + /** The boxes one host is responsible for, as a subquery to scope a write. */ + function boxesOn(tx: Parameters[0]>[0], machineId: string) { + return tx + .select({ id: BoxTable.id }) + .from(BoxTable) + .where(and(eq(BoxTable.machineId, machineId), isNull(BoxTable.timeDeleted))); + } + + /** + * Move a run from one exact state to another, or do nothing at all. + * + * This is the claim, and it is why `setState` is not enough on its own: + * updating on the id alone means two agents polling the same work both + * succeed and both start the same box. The current state is part of the + * `where` clause, so the database decides the winner and the loser gets + * null rather than a row. There is one host today, which is exactly why + * this would otherwise be built wrong and stay wrong. + * + * The host is in the same `where` clause. The caller checking first is not + * the same thing as the write being scoped, and only one of the two is + * still true when somebody adds a second caller. + */ + export const compareAndSetState = fn( + z.object({ + id: Info.shape.id, + machineId: z.string(), + from: z.enum(SessionState.enumValues), + to: z.enum(SessionState.enumValues), + errorMessage: Info.shape.errorMessage + }), + async (input) => { + const now = sql`now()`; + return Database.use(async (tx) => { + return tx + .update(SessionTable) + .set({ + state: input.to, + errorMessage: input.to === 'failed' ? (input.errorMessage ?? null) : null, + ...(input.to === 'live' + ? { timeStarted: sql`coalesce(${SessionTable.timeStarted}, ${now})` } + : {}), + ...(input.to === 'ended' || input.to === 'failed' + ? { timeStopped: sql`coalesce(${SessionTable.timeStopped}, ${now})` } + : {}) + }) + .where( + and( + eq(SessionTable.id, input.id), + eq(SessionTable.state, input.from), + isNull(SessionTable.timeDeleted), + inArray(SessionTable.boxId, boxesOn(tx, input.machineId)) + ) + ) + .returning() + .then((rows) => { + const row = rows.at(0); + return row ? serialize(row) : null; + }); + }); + } + ); + + /** + * What happened when an agent reported a state. + * + * Four outcomes that look alike from a distance and are not, which is the + * whole reason this is not a boolean: + * + * - `forbidden` — no such run, or it is not on this host. One answer for + * both, so reporting states at ids cannot be used to discover them. + * - `unchanged` — already in that state. A retry after a lost response is + * not a broken agent and must not be told it is. + * - `illegal` — not a transition that exists. The row does not move. + * - `lost` — a legal transition that something else got to first. + * - `moved` — it happened. + */ + export type TransitionOutcome = 'forbidden' | 'unchanged' | 'illegal' | 'lost' | 'moved'; + + export interface TransitionResult { + outcome: TransitionOutcome; + session: Info | null; + } + + export const transition = fn( + z.object({ + id: Info.shape.id, + machineId: z.string(), + state: z.enum(SessionState.enumValues), + errorMessage: Info.shape.errorMessage + }), + async (input): Promise => { + const current = await forMachine({ id: input.id, machineId: input.machineId }); + if (!current) return { outcome: 'forbidden', session: null }; + if (current.state === input.state) return { outcome: 'unchanged', session: current }; + if (!NEXT_STATES[current.state].includes(input.state)) { + return { outcome: 'illegal', session: current }; + } + + const moved = await compareAndSetState({ + id: input.id, + machineId: input.machineId, + from: current.state, + to: input.state, + errorMessage: input.errorMessage + }); + // The state read above is not the state written below, and the gap + // is where two agents race. Nothing moved means somebody else did. + if (!moved) return { outcome: 'lost', session: current }; + return { outcome: 'moved', session: moved }; + } + ); + + export interface TicketResult { + outcome: 'forbidden' | 'closed' | 'published'; + session: Info | null; + } + + /** + * Publish a ticket for a run, on behalf of the host it is placed on. + * + * A ticket may appear while the state is still `starting` — it is + * republished as addresses are discovered, so the client polls and re-reads + * rather than keeping the first one. A run that has stopped is refused: an + * address for something that is not there can only mislead whoever is + * still polling. + */ + export const publishTicket = fn( + z.object({ + id: Info.shape.id, + machineId: z.string(), + ticket: z.string().min(1) + }), + async (input): Promise => { + const current = await forMachine({ id: input.id, machineId: input.machineId }); + if (!current) return { outcome: 'forbidden', session: null }; + + return Database.use(async (tx) => { + return tx + .update(SessionTable) + .set({ ticket: input.ticket }) + .where( + and( + eq(SessionTable.id, input.id), + notInArray(SessionTable.state, ['ended', 'failed']), + isNull(SessionTable.timeDeleted), + inArray(SessionTable.boxId, boxesOn(tx, input.machineId)) + ) + ) + .returning() + .then((rows): TicketResult => { + const row = rows.at(0); + return row + ? { outcome: 'published', session: serialize(row) } + : { outcome: 'closed', session: current }; + }); + }); + } + ); + export function serialize(input: typeof SessionTable.$inferSelect): z.infer { return { id: input.id, diff --git a/packages/core/src/session/session.test.ts b/packages/core/src/session/session.test.ts index 12206713..d2803a63 100644 --- a/packages/core/src/session/session.test.ts +++ b/packages/core/src/session/session.test.ts @@ -41,7 +41,7 @@ async function scene(label: string, steamAppId: number) { label, tier: 'sm' }); - return { owner, box, gameId: await newGame(steamAppId) }; + return { owner, machineId, box, gameId: await newGame(steamAppId) }; } afterAll(async () => { @@ -173,3 +173,244 @@ describe('Session', () => { expect(await Session.listByBox(box.id)).toHaveLength(0); }); }); + +describe('Session jobs', () => { + test('a requested session is the job, and it carries its kind', async () => { + const { owner, machineId, box, gameId } = await scene('ses-job-kind', 5410); + const session = await Session.create({ + id: Identifier.ascending('session'), + boxId: box.id, + gameId, + linkedAccountId: owner.linkedAccountId + }); + + const jobs = await Session.listJobsForMachine(machineId); + expect(jobs).toHaveLength(1); + // The kind is on the wire from the first day there is only one, so the + // second kind is an addition rather than a redesign. + expect(jobs[0]!.kind).toBe('session.start'); + expect(jobs[0]!.sessionId).toBe(session.id); + expect(jobs[0]!.boxId).toBe(box.id); + expect(jobs[0]!.boxTier).toBe('sm'); + expect(jobs[0]!.gameId).toBe(gameId); + expect(jobs[0]!.steamAppId).toBe(5410); + expect(jobs[0]!.linkedAccountId).toBe(owner.linkedAccountId); + }); + + test('a job belongs to the machine its box is placed on and to no other', async () => { + const mine = await scene('ses-job-mine', 5411); + const theirs = await scene('ses-job-theirs', 5412); + + const session = await Session.create({ + id: Identifier.ascending('session'), + boxId: theirs.box.id, + gameId: theirs.gameId, + linkedAccountId: theirs.owner.linkedAccountId + }); + + // The scope is the join, not a filter the caller asks for. A machine + // credential is a long-lived secret on hardware in somebody's home, so + // what one leaking can reach is decided here. + expect(await Session.listJobsForMachine(mine.machineId)).toHaveLength(0); + expect((await Session.listJobsForMachine(theirs.machineId)).map((j) => j.sessionId)).toEqual([ + session.id + ]); + }); + + test('only a requested session is work; a claimed one is not offered again', async () => { + const { owner, machineId, box, gameId } = await scene('ses-job-claimed', 5413); + const session = await Session.create({ + id: Identifier.ascending('session'), + boxId: box.id, + gameId, + linkedAccountId: owner.linkedAccountId + }); + + expect(await Session.listJobsForMachine(machineId)).toHaveLength(1); + await Session.transition({ + id: session.id, + machineId, + state: 'starting', + errorMessage: null + }); + expect(await Session.listJobsForMachine(machineId)).toHaveLength(0); + }); +}); + +describe('Session claim', () => { + async function requested(label: string, steamAppId: number) { + const s = await scene(label, steamAppId); + const session = await Session.create({ + id: Identifier.ascending('session'), + boxId: s.box.id, + gameId: s.gameId, + linkedAccountId: s.owner.linkedAccountId + }); + return { ...s, session }; + } + + test('the claim is a compare-and-set, so the second attempt finds nothing to move', async () => { + const { machineId, session } = await requested('ses-cas', 5420); + + const won = await Session.compareAndSetState({ + id: session.id, + machineId, + from: 'requested', + to: 'starting', + errorMessage: null + }); + expect(won?.state).toBe('starting'); + + // The same attempt again. The row is no longer `requested`, so the + // update matches nothing — which is what stops two agents from both + // starting the same box. Updating on the id alone would succeed twice. + const lost = await Session.compareAndSetState({ + id: session.id, + machineId, + from: 'requested', + to: 'starting', + errorMessage: null + }); + expect(lost).toBeNull(); + expect((await Session.fromID(session.id))?.state).toBe('starting'); + }); + + test('a machine that is not the box’s host cannot move the row', async () => { + const { session } = await requested('ses-cas-mine', 5421); + const other = await scene('ses-cas-other', 5422); + + const result = await Session.transition({ + id: session.id, + machineId: other.machineId, + state: 'starting', + errorMessage: null + }); + expect(result.outcome).toBe('forbidden'); + expect((await Session.fromID(session.id))?.state).toBe('requested'); + + // And the compare-and-set is scoped in the same query, not only by the + // classification above it. + expect( + await Session.compareAndSetState({ + id: session.id, + machineId: other.machineId, + from: 'requested', + to: 'starting', + errorMessage: null + }) + ).toBeNull(); + }); + + test('a session that does not exist is refused the same way as one that is not yours', async () => { + const other = await scene('ses-cas-ghost', 5423); + const result = await Session.transition({ + // A well-formed id for a row that was never written. + id: Identifier.ascending('session'), + machineId: other.machineId, + state: 'starting', + errorMessage: null + }); + // Same answer as somebody else's session: an agent must not be able to + // learn which ids exist by reporting states at them. + expect(result.outcome).toBe('forbidden'); + }); + + test('re-reporting the state you already reported changes nothing', async () => { + const { machineId, session } = await requested('ses-repeat', 5424); + + await Session.transition({ id: session.id, machineId, state: 'starting', errorMessage: null }); + const again = await Session.transition({ + id: session.id, + machineId, + state: 'starting', + errorMessage: null + }); + // A retry after a lost response is not a broken agent. + expect(again.outcome).toBe('unchanged'); + expect(again.session?.state).toBe('starting'); + }); + + test('a transition off the table is refused and the row does not move', async () => { + const { machineId, session } = await requested('ses-illegal', 5425); + + // Skipping `starting` means nothing ever holds the claim, and the claim + // is the only mutual exclusion here — so it is refused however tempting + // the shortcut looks. + const skipped = await Session.transition({ + id: session.id, + machineId, + state: 'live', + errorMessage: null + }); + expect(skipped.outcome).toBe('illegal'); + expect((await Session.fromID(session.id))?.state).toBe('requested'); + + await Session.transition({ id: session.id, machineId, state: 'starting', errorMessage: null }); + await Session.transition({ id: session.id, machineId, state: 'failed', errorMessage: 'no' }); + + // Terminal is terminal: a dead session cannot be resurrected. + const raised = await Session.transition({ + id: session.id, + machineId, + state: 'live', + errorMessage: null + }); + expect(raised.outcome).toBe('illegal'); + expect((await Session.fromID(session.id))?.state).toBe('failed'); + }); + + test('the timestamps survive a duplicate report, which is what billing rests on', async () => { + const { machineId, session } = await requested('ses-idempotent', 5426); + await Session.transition({ id: session.id, machineId, state: 'starting', errorMessage: null }); + const live = await Session.transition({ + id: session.id, + machineId, + state: 'live', + errorMessage: null + }); + expect(live.session?.timeStarted).not.toBeNull(); + + const repeat = await Session.transition({ + id: session.id, + machineId, + state: 'live', + errorMessage: null + }); + expect(repeat.session?.timeStarted).toBe(live.session!.timeStarted); + }); + + test('publishing a ticket is scoped to the host too', async () => { + const { machineId, session } = await requested('ses-ticket-scope', 5427); + const other = await scene('ses-ticket-other', 5428); + + await Session.transition({ id: session.id, machineId, state: 'starting', errorMessage: null }); + + const refused = await Session.publishTicket({ + id: session.id, + machineId: other.machineId, + ticket: 'stolen' + }); + expect(refused.outcome).toBe('forbidden'); + expect((await Session.fromID(session.id))?.ticket).toBeNull(); + + // A ticket may appear while the state is still `starting`. + const first = await Session.publishTicket({ id: session.id, machineId, ticket: 'one' }); + expect(first.outcome).toBe('published'); + expect(first.session?.ticket).toBe('one'); + expect(first.session?.state).toBe('starting'); + + const second = await Session.publishTicket({ id: session.id, machineId, ticket: 'two' }); + expect(second.session?.ticket).toBe('two'); + }); + + test('a stopped session has no address to publish', async () => { + const { machineId, session } = await requested('ses-ticket-dead', 5429); + await Session.transition({ id: session.id, machineId, state: 'starting', errorMessage: null }); + await Session.transition({ id: session.id, machineId, state: 'live', errorMessage: null }); + await Session.transition({ id: session.id, machineId, state: 'ended', errorMessage: null }); + + const result = await Session.publishTicket({ id: session.id, machineId, ticket: 'late' }); + expect(result.outcome).toBe('closed'); + expect((await Session.fromID(session.id))?.ticket).toBeNull(); + }); +});