Rebuild the shell, add Calendar and Learn, and govern reads
CI / verify (push) Successful in 3m45s
CI / publish (push) Has been skipped

Seven parallel agents and an adversarial verification pass. The three things
worth knowing before reading the diff:

RBAC WAS ALREADY BUILT. docs/build-plan.md marks F2 and F3 outstanding and is
stale — packages/core/src/permissions.ts and lib/mutation.ts shipped long ago.
So this does not rebuild them; it closes the gaps an audit found. The big one
is that reads were entirely ungoverned: every GET was "any authenticated
member", so a junior demand rep and a research contractor could both pull
per-block supplier cost and break-even prices from /api/capacity/margin, and
every contract's negotiated terms. For a company whose margin is the business,
that was the hole that mattered. Adds book:read / economics:read / team:read,
a readGuard middleware, and a `viewer` role below member.

THE BUTTON AND THE 403 DISAGREED — the exact thing F3 said must never happen.
Contracts.tsx never called can() at all, so its save button was always enabled
against a server requiring contract:sign; Capacity.tsx gated commitment
creation on deal:write/demand while the server wanted commitment:write/supply.

POST /api/activities was the one write bypassing executeMutation: no capability
check, and any member could mutate accounts.lastActivityAt as a side effect.
It is now a proper mutation() behind activity:write.

The shell becomes three panes — a collapsible shadcn sidebar with an account
switcher on the Piggy accent, a header with real search, and Piggy docked to
the right, page-aware and persistent across navigation. The phone keeps its
bottom tab bar, which is the thing this product already beat trycompai/crm on,
and gains the sidebar as a sheet.

Calendar is a projection over thirteen dated sources rather than a new table,
because a table would duplicate dates that already live on contracts, deals and
commitments and would drift — and one ledger answering the question is the
whole argument. It surfaces export_authorizations and compliance_artifacts,
which had indexed expires_at columns, schema comments saying they must be
alerted on, and no read endpoint or UI anywhere.

Learn carries two tracks. Concepts are members-only; the platform track can be
opened with a share code by someone with no account. The code mints a scoped
learn-only token and never a Principal — every route here resolves a principal
and then checks capabilities, so a principal-minting code would be one missing
check away from leaking the book. "Only platform-track rows may be code-visible"
is a database CHECK constraint as well as a write-path rule, and a test asserts
a valid learn token still gets 401 on /api/dashboard, /api/accounts and
/api/contracts — the same invariant scripts/deploy.sh refuses to ship without.

CD becomes tag-to-ship. CI publishes an image to the Gitea registry on a
release-* tag and cloud-2 pulls it, so no credential on the shared runner can
execute anything on production — by construction rather than by policy. Both
halves of deploy.sh's original rule survive: nothing on the runner reaches the
host, and a human still decides when it ships. deploy.sh gains a rollback and a
public-origin check, and PIG_IMAGE now reaches compose through `sudo env`,
without which sudo's env_reset silently resolved every release to pig:local.

Tests 141 -> 261.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
This commit is contained in:
2026-08-13 15:02:48 -07:00
parent 6cf80747cc
commit 13dec6b4b8
102 changed files with 28638 additions and 913 deletions
+34 -50
View File
@@ -26,7 +26,6 @@ import {
} from '@pig/db';
import {
ACCENTS,
ACTIVITY_TYPES,
DEMAND_STAGES,
SECURITY_TIERS,
SUPPLY_STAGES,
@@ -58,13 +57,17 @@ import { createRecordRoutes } from './routes/records';
import { createImportRoutes } from './routes/imports';
import { createGoogleSheetsRoutes } from './routes/google-sheets';
import { createContractRoutes } from './routes/contracts';
import { createPiggyChatRoutes } from './routes/piggy-chat';
import { createPiggyChatRoutes, platformPiggyEnabled } from './routes/piggy-chat';
import { createAdminSettingsRoutes } from './routes/admin-settings';
import { createSlackRoutes, SLACK_CAPACITY_COMMAND_PATH } from './routes/slack';
import { createBuzzRoutes } from './routes/buzz';
import { createIntegrationSettingsRoutes } from './routes/integration-settings';
import { createNotionImportRoutes, NOTION_OAUTH_CALLBACK_PATH } from './routes/notion-import';
import { createGrowthRoutes } from './routes/growth';
import { createCalendarRoutes } from './routes/calendar';
import { createLearnRoutes, LEARN_ACCESS_PATH, LEARN_PUBLIC_PATH } from './routes/learn';
import { createReadGuardRoutes } from './routes/read-guards';
import { createActivityRoutes } from './routes/activities';
import { NotificationOutbox } from './services/notification-outbox';
type Env = { Variables: { principal: Principal } };
@@ -142,6 +145,13 @@ export function createApp(
path === '/api/register'
|| path === SLACK_CAPACITY_COMMAND_PATH
|| path === NOTION_OAUTH_CALLBACK_PATH
// Learn is reachable with a share code and no account. These two paths
// are exact-string matches, deliberately: /api/learn and
// /api/learn/resources/* stay behind the authenticator, and the public
// reader is structurally incapable of naming a row that is not both
// platform-track and code-visible.
|| path === LEARN_ACCESS_PATH
|| path === LEARN_PUBLIC_PATH
) {
return next();
}
@@ -156,6 +166,19 @@ export function createApp(
return next();
});
/*
* Read authorisation, mounted before every handler it guards.
*
* Hono runs matched handlers in registration order, so a guard registered
* after its route never runs and returns 200 while looking correct. That is
* why this sits here rather than beside the feature routes below, and why
* read-governance.test.ts pins the ordering in both directions.
*
* The policy is one table in read-guards.ts precisely so that "who can see
* cost?" has a single answer rather than one per route.
*/
app.route('/', createReadGuardRoutes());
// ---------------------------------------------------------------- identity
app.get('/api/me', (c) => {
@@ -220,12 +243,20 @@ export function createApp(
}));
app.route('/', createContractRoutes(db));
app.route('/', createGrowthRoutes(db));
app.route('/', createCalendarRoutes(db));
app.route('/', createLearnRoutes(db));
app.route(
'/',
createPiggyChatRoutes({
enabled: config.PIGGY_ENABLED,
internalUrl: config.PIGGY_INTERNAL_URL,
internalToken: config.PIGGY_INTERNAL_TOKEN,
// Without this the stored toggle is never consulted and isAvailable()
// short-circuits to the environment variable, which is the bug the
// resolver exists to fix. The tests inject their own resolver, so they
// stay green whether or not this line is here — it is the composition
// that has to be right.
resolvePiggyEnabled: platformPiggyEnabled(config, db),
}),
);
app.route('/', createSlackRoutes(config, db, capacity));
@@ -358,54 +389,7 @@ export function createApp(
app.route('/', createCapacityWriteRoutes(db));
app.route('/', createFactsRoute(db));
// ------------------------------------------------------------- activities
const activitySchema = z.object({
accountId: z.string().uuid().optional(),
contactId: z.string().uuid().optional(),
demandDealId: z.string().uuid().optional(),
supplyDealId: z.string().uuid().optional(),
type: z.enum(ACTIVITY_TYPES),
subject: z.string().min(1).max(200),
body: z.string().max(8000).optional(),
occurredAt: z.string().datetime().optional(),
externalId: z.string().max(200).optional(),
});
app.post('/api/activities', async (c) => {
const p = c.get('principal');
const parsed = activitySchema.safeParse(await c.req.json());
if (!parsed.success) {
return c.json({ error: 'Invalid activity', issues: parsed.error.issues }, 400);
}
const { occurredAt, ...rest } = parsed.data;
const when = occurredAt ? new Date(occurredAt) : new Date();
const [created] = await db
.insert(activities)
.values({
...rest,
occurredAt: when,
actorUserId: p.userId,
// An agent acting for someone is recorded as such, so the log
// distinguishes what a person did from what was done on their behalf.
actorAgent: p.via === 'api_key' ? 'agent' : null,
source: p.via === 'api_key' ? 'agent' : 'manual',
})
// An `externalId` collision means this event was already synced from
// Slack or Buzz; silently ignoring the duplicate keeps sync idempotent.
.onConflictDoNothing()
.returning();
if (rest.accountId) {
await db
.update(accounts)
.set({ lastActivityAt: when })
.where(eq(accounts.id, rest.accountId));
}
return c.json(created ?? { deduplicated: true }, created ? 201 : 200);
});
app.route('/', createActivityRoutes(db));
// ---------------------------------------------------------------- capacity
+67 -10
View File
@@ -24,12 +24,17 @@ import type { Database } from '@pig/db';
import { apiKeys, teamMemberships, users } from '@pig/db';
import {
permissionGranted,
resolvePermissionGrants,
type Capability,
resolveReadPermissionGrants,
resolveWritePermissionGrants,
roleMeets,
TEAM_CAPABILITY_RULES,
type GlobalCapability,
type PermissionGrant,
type ReadCapability,
type Team,
type TeamCapability,
type TeamRole,
type WriteCapability,
} from '@pig/core';
import { createHash, timingSafeEqual } from 'node:crypto';
import type { Config } from './config';
@@ -231,8 +236,7 @@ export function hasTeamAccess(
if (principal.isPlatformAdmin) return true;
const membership = principal.teams.find((t) => t.team === team);
if (!membership) return false;
const rank: Record<TeamRole, number> = { member: 0, lead: 1, admin: 2 };
return rank[membership.role] >= rank[minimumRole];
return roleMeets(membership.role, minimumRole);
}
export function requireScope(principal: Principal, scope: string): void {
@@ -240,13 +244,22 @@ export function requireScope(principal: Principal, scope: string): void {
throw new AuthError(`This credential lacks the '${scope}' scope.`, 403, 'insufficient_scope');
}
/** Effective grants include credential scope, not merely the owner's roles. */
/**
* Effective grants include credential scope, not merely the owner's roles.
*
* Read and write scopes are filtered separately. Before read capabilities
* existed a read-only key resolved to no grants at all, which was right then
* and would now be wrong: it would tell `/api/me` that a read-only agent may
* not read, and the browser would grey out a page the server happily serves.
*/
export function effectivePermissions(principal: Principal): PermissionGrant[] {
if (!principal.scopes.includes('write')) return [];
return resolvePermissionGrants(principal);
const grants: PermissionGrant[] = [];
if (principal.scopes.includes('read')) grants.push(...resolveReadPermissionGrants(principal));
if (principal.scopes.includes('write')) grants.push(...resolveWritePermissionGrants(principal));
return grants;
}
export function requireCapability(principal: Principal, capability: Capability): void;
export function requireCapability(principal: Principal, capability: GlobalCapability): void;
export function requireCapability(
principal: Principal,
capability: TeamCapability,
@@ -254,14 +267,58 @@ export function requireCapability(
): void;
export function requireCapability(
principal: Principal,
capability: Capability,
capability: WriteCapability,
team?: Team,
): void {
requireScope(principal, 'write');
if (permissionGranted(resolvePermissionGrants(principal), capability, team)) return;
if (permissionGranted(resolveWritePermissionGrants(principal), capability, team)) return;
throw new AuthError(
`This principal lacks the '${capability}' capability${team ? ` for ${team}` : ''}.`,
403,
'insufficient_permission',
);
}
/**
* "May they do this on *some* team?"
*
* Separate from `requireCapability` and deliberately harder to type by
* accident. Passing no team to the old `requireCapability` silently meant this
* — which is how a research-team admin could bulk-import demand deals — so the
* overloads above now refuse it and every remaining any-team check has to say
* so in its own name. Use it only where no team is knowable yet: listing the
* spreadsheets in someone's Drive, before an entity has been chosen. The
* moment the target is known, go back to `requireCapability` with its team.
*/
export function requireAnyTeamCapability(
principal: Principal,
capability: TeamCapability,
): void {
requireScope(principal, 'write');
const grants = resolveWritePermissionGrants(principal);
for (const team of TEAM_CAPABILITY_RULES[capability].teams) {
if (permissionGranted(grants, capability, team)) return;
}
throw new AuthError(
`This principal lacks the '${capability}' capability on any team.`,
403,
'insufficient_permission',
);
}
/**
* Reads are governed too. The 'read' scope is checked rather than 'write'
* because a read-only API key is exactly the credential this must admit.
*/
export function requireReadCapability(
principal: Principal,
capability: ReadCapability,
): void {
requireScope(principal, 'read');
if (permissionGranted(resolveReadPermissionGrants(principal), capability)) return;
throw new AuthError(
`This principal lacks the '${capability}' capability.`,
403,
'insufficient_permission',
);
}
+27 -11
View File
@@ -1,4 +1,5 @@
import type { ActivityType, GlobalCapability, Team, TeamCapability } from '@pig/core';
import { isTeamCapability } from '@pig/core';
import type { Database } from '@pig/db';
import { activities } from '@pig/db';
import type { Context, Handler } from 'hono';
@@ -55,9 +56,18 @@ export interface MutationActivity {
meta?: Record<string, unknown>;
}
/**
* `'self'` is for the one write whose own row IS the audit event: logging an
* activity. Inserting an audit row about it would double every synced call in
* the feed. It is a literal rather than an omitted field so that audit can
* never be skipped by forgetting to write one — the type still demands an
* answer, and `'self'` is a visible, greppable claim.
*/
export type MutationAudit = MutationActivity | 'self';
export interface MutationResult<Result> {
data: Result;
activity: MutationActivity;
activity: MutationAudit;
}
interface MutationContext<Input> {
@@ -80,11 +90,15 @@ function enforcePermission(principal: Principal, permission: PermissionRequireme
permission.authorize(principal);
return;
}
if (permission.capability === 'settings:admin') {
requireCapability(principal, permission.capability);
// Discriminated by the capability itself rather than by a hard-coded
// 'settings:admin' check, which quietly sent any future global capability
// down the team-scoped branch with an undefined team — the "passes on any
// team" bug, reintroduced by omission.
if (isTeamCapability(permission.capability)) {
requireCapability(principal, permission.capability, permission.team as Team);
return;
}
requireCapability(principal, permission.capability, permission.team);
requireCapability(principal, permission.capability as GlobalCapability);
}
/**
@@ -131,13 +145,15 @@ export async function executeMutation<Schema extends ZodTypeAny, Result>(
};
const result = await definition.mutate(context);
await tx.insert(activities).values({
...result.activity,
actorUserId: principal.userId,
actorAgent: principal.via === 'api_key' ? 'agent' : null,
source: principal.via === 'api_key' ? 'agent' : 'manual',
occurredAt: now,
});
if (result.activity !== 'self') {
await tx.insert(activities).values({
...result.activity,
actorUserId: principal.userId,
actorAgent: principal.via === 'api_key' ? 'agent' : null,
source: principal.via === 'api_key' ? 'agent' : 'manual',
occurredAt: now,
});
}
return result.data;
});
}
+30
View File
@@ -0,0 +1,30 @@
/**
* Authorisation for reads.
*
* The write path has had one chokepoint since F2 — `executeMutation` — and
* reads had none. Every GET was "any authenticated member", so a research
* contractor and a demand lead saw supplier cost per GPU-hour, break-even
* price and the full negotiated terms of every contract identically. For a
* company whose margin is the product, that was the hole that mattered.
*
* This is the reading half of the same chokepoint. It is thin on purpose:
* capability in, middleware out, and the AuthError it throws is mapped to HTTP
* by `app.onError` exactly as the write path's is, so a read denial and a write
* denial are indistinguishable in shape to a client.
*
* `growth.ts` had the shape of this already but keyed on API-key *scope*, which
* answers "is this credential allowed to read anything?" and not "is this
* person allowed to read *this*". Scope is a property of the credential; the
* capability is a property of the person. Both are checked here.
*/
import type { ReadCapability } from '@pig/core';
import type { MiddlewareHandler } from 'hono';
import { requireReadCapability } from './auth';
import type { ApiEnv } from './mutation';
export function readGuard(capability: ReadCapability): MiddlewareHandler<ApiEnv> {
return async (context, next) => {
requireReadCapability(context.get('principal'), capability);
await next();
};
}
+136
View File
@@ -0,0 +1,136 @@
/**
* Logging an activity.
*
* This was the one write in PIG that never went through `executeMutation`: it
* lived inline in `app.ts`, checked no capability at all, and would insert an
* activity against any `accountId` a caller cared to name — then move that
* account's `lastActivityAt`, which is what the account list sorts on. Any
* member, and any write-scoped API key, could therefore reorder somebody
* else's book and plant a fabricated call in the audit trail of an account
* they have no relationship with.
*
* Two things are checked, in two places, deliberately:
*
* 1. Up front, before the body is read: does this principal hold
* `activity:write` on *any* team? A caller with none must not get to probe
* validation rules for a write they can never perform.
* 2. Inside the transaction, once the referenced account has been read: do
* they hold it on a team that account is actually on? The side is a
* property of the row, so it cannot be known before the row is fetched.
*
* That is the same shape as `ensureSidePermission` in contracts.ts, and for the
* same reason.
*/
import { ACTIVITY_TYPES, type AccountSide, type Team } from '@pig/core';
import type { Activity, Database } from '@pig/db';
import { accounts, activities } from '@pig/db';
import { eq } from 'drizzle-orm';
import { Hono } from 'hono';
import { z } from 'zod';
import { AuthError, requireAnyTeamCapability, requireCapability, type Principal } from '../lib/auth';
import type { ApiEnv, MutationDefinition } from '../lib/mutation';
import { MutationError, mutation } from '../lib/mutation';
const activitySchema = z
.object({
accountId: z.string().uuid().optional(),
contactId: z.string().uuid().optional(),
demandDealId: z.string().uuid().optional(),
supplyDealId: z.string().uuid().optional(),
type: z.enum(ACTIVITY_TYPES),
subject: z.string().min(1).max(200),
body: z.string().max(8000).optional(),
occurredAt: z.string().datetime().optional(),
externalId: z.string().max(200).optional(),
})
.strict();
export interface LoggedActivity {
activity: Activity | null;
/** True when an `externalId` collision meant the event was already synced. */
deduplicated: boolean;
}
/**
* Research consumes capacity but keeps no commercial book, so an activity is
* always a supply-side or demand-side event. A `both` account admits either.
*/
function requireSidePermission(principal: Principal, side: AccountSide): void {
const sides: Team[] = side === 'both' ? ['supply', 'demand'] : [side];
for (const team of sides) {
try {
requireCapability(principal, 'activity:write', team);
return;
} catch (error) {
if (!(error instanceof AuthError)) throw error;
}
}
throw new AuthError(
`This principal cannot log activity against a ${side}-side account.`,
403,
'insufficient_permission',
);
}
export function createActivityMutationDefinition(): MutationDefinition<
typeof activitySchema,
LoggedActivity
> {
return {
schema: activitySchema,
permission: { authorize: (principal) => requireAnyTeamCapability(principal, 'activity:write') },
invalidMessage: 'Invalid activity.',
async mutate({ input, principal, tx, now }) {
const { occurredAt, accountId, ...rest } = input;
// A backdated entry is the normal case for sync, so the caller's
// timestamp wins over `now` — unlike the audit rows this convention
// usually writes, where `now` is the point.
const when = occurredAt ? new Date(occurredAt) : now;
if (accountId) {
const [account] = await tx
.select({ side: accounts.side })
.from(accounts)
.where(eq(accounts.id, accountId))
.limit(1);
if (!account) throw MutationError.notFound('Account');
requireSidePermission(principal, account.side as AccountSide);
}
const [created] = await tx
.insert(activities)
.values({
...rest,
accountId,
occurredAt: when,
actorUserId: principal.userId,
// An agent acting for someone is recorded as such, so the log
// distinguishes what a person did from what was done on their behalf.
actorAgent: principal.via === 'api_key' ? 'agent' : null,
source: principal.via === 'api_key' ? 'agent' : 'manual',
})
// An `externalId` collision means this event was already synced from
// Slack or Buzz; silently ignoring the duplicate keeps sync idempotent.
.onConflictDoNothing()
.returning();
// Only on a real insert. Bumping it on a deduplicated replay would let a
// repeated sync keep an account at the top of the list forever.
if (created && accountId) {
await tx.update(accounts).set({ lastActivityAt: when }).where(eq(accounts.id, accountId));
}
return {
data: { activity: created ?? null, deduplicated: !created },
// The inserted row is the audit event. See `MutationAudit`.
activity: 'self',
};
},
};
}
export function createActivityRoutes(db: Database): Hono<ApiEnv> {
const routes = new Hono<ApiEnv>();
routes.post('/api/activities', mutation(db, createActivityMutationDefinition()));
return routes;
}
+494
View File
@@ -0,0 +1,494 @@
/**
* GET /api/calendar and the CRUD for the one table it owns.
*
* The read endpoint is the first in this API to accept a date range and filter
* on it server-side. Every other list route is `order by updated_at desc limit
* 300` with the browser filtering afterwards, which means the records dated
* inside a quarter are not guaranteed to be in the response — the failure this
* route exists to remove. Nothing here computes anything; the projection lives
* in the service, per the rule that intelligence never lives in the API.
*/
import { and, eq } from 'drizzle-orm';
import { Hono, type Context } from 'hono';
import { z } from 'zod';
import {
CALENDAR_ENTRY_KINDS,
CALENDAR_EVENT_KINDS,
isCalendarEventKind,
isValidTimeZone,
parseQuarter,
permissionGranted,
quarterBounds,
quarterBoundsFor,
type CalendarEventKind,
} from '@pig/core';
import { accounts, calendarEntries, demandDeals, supplyDeals, users } from '@pig/db';
import type { CalendarEntry, Database } from '@pig/db';
import { effectivePermissions, requireCapability, type Principal } from '../lib/auth';
import {
MutationError,
apiError,
bodylessMutation,
mutation,
type ApiEnv,
type MutationDefinition,
} from '../lib/mutation';
import { CalendarService } from '../services/calendar';
/**
* A dated item belongs to whoever runs the motion, so either pipeline's
* write-capable members may keep the calendar. Mirrors the treatment contracts
* already give a capability that is meaningful on both sides.
*/
export function requireCalendarWrite(principal: Principal): void {
const grants = effectivePermissions(principal);
if (
permissionGranted(grants, 'deal:write', 'supply') ||
permissionGranted(grants, 'deal:write', 'demand')
) {
return;
}
// Re-run the check so the caller gets the standard 403 envelope rather than
// a bespoke one, and so a scopeless credential is reported as such.
requireCapability(principal, 'deal:write', 'demand');
}
export function calendarReadAllowed(scopes: readonly string[]): boolean {
return scopes.includes('read');
}
/**
* Comma-separated, and an unknown kind is an error rather than a silent empty
* result — a typo in `kinds` that returns nothing looks exactly like a quiet
* quarter.
*/
export function parseKinds(raw: string | undefined): CalendarEventKind[] | undefined {
if (!raw) return undefined;
const requested = raw
.split(',')
.map((kind) => kind.trim())
.filter(Boolean);
if (!requested.length) return undefined;
const unknown = requested.filter((kind) => !isCalendarEventKind(kind));
if (unknown.length) {
throw new MutationError(
'invalid_kinds',
`Unknown calendar event kind(s): ${unknown.join(', ')}. Known kinds: ${CALENDAR_EVENT_KINDS.join(', ')}.`,
400,
);
}
return requested as CalendarEventKind[];
}
const isoDate = z.string().datetime();
const nullableId = z.string().uuid().nullable();
const entryFields = {
title: z.string().min(1).max(240).optional(),
description: z.string().max(8000).nullable().optional(),
kind: z.enum(CALENDAR_ENTRY_KINDS).optional(),
startsAt: isoDate.optional(),
endsAt: isoDate.nullable().optional(),
allDay: z.boolean().optional(),
ownerUserId: nullableId.optional(),
accountId: nullableId.optional(),
demandDealId: nullableId.optional(),
supplyDealId: nullableId.optional(),
completedAt: isoDate.nullable().optional(),
};
const createEntrySchema = z.object({
...entryFields,
title: z.string().min(1).max(240),
startsAt: isoDate,
});
const updateEntrySchema = z.object(entryFields);
function date(value: string | null | undefined): Date | null | undefined {
return value === undefined ? undefined : value === null ? null : new Date(value);
}
function requiredRouteParam(params: Readonly<Record<string, string>>): string {
const value = params.id;
if (!value) {
throw new MutationError('invalid_route_parameter', "Route parameter 'id' is required.", 400);
}
return value;
}
function writtenRow<Row>(row: Row | undefined): Row {
if (row === undefined) {
throw new Error('Calendar entry write completed without returning a row.');
}
return row;
}
type Transaction = Parameters<Parameters<Database['transaction']>[0]>[0];
/**
* The nullable foreign keys are polymorphic, so nothing in the schema stops an
* entry pointing at an account and a deal belonging to someone else. Checked
* here, where the intent is known.
*/
async function requireRelationships(
tx: Transaction,
record: {
/**
* Only ever the client's own choice. The default — the author's own id —
* is a user we have just authenticated, so re-reading it would be a query
* per create to confirm something the request already proved.
*/
ownerUserId?: string | null;
accountId?: string | null;
demandDealId?: string | null;
supplyDealId?: string | null;
startsAt: Date;
endsAt?: Date | null;
},
): Promise<void> {
if (record.endsAt && record.endsAt < record.startsAt) {
throw new MutationError('invalid_window', 'An entry cannot end before it starts.', 400);
}
if (record.ownerUserId) {
// Assigning to someone who has since been removed is an ordinary client
// mistake, and without this it surfaces as a 500 from the foreign key
// rather than the 404 every other polymorphic reference here returns.
const [owner] = await tx
.select({ id: users.id })
.from(users)
.where(eq(users.id, record.ownerUserId))
.limit(1);
if (!owner) throw MutationError.notFound('User');
}
if (record.accountId) {
const [account] = await tx
.select({ id: accounts.id })
.from(accounts)
.where(eq(accounts.id, record.accountId))
.limit(1);
if (!account) throw MutationError.notFound('Account');
}
if (record.demandDealId) {
const [deal] = await tx
.select({ accountId: demandDeals.accountId })
.from(demandDeals)
.where(eq(demandDeals.id, record.demandDealId))
.limit(1);
if (!deal) throw MutationError.notFound('Demand deal');
if (record.accountId && deal.accountId !== record.accountId) {
throw new MutationError(
'relationship_mismatch',
'Demand deal belongs to a different account.',
409,
);
}
}
if (record.supplyDealId) {
const [deal] = await tx
.select({ accountId: supplyDeals.accountId })
.from(supplyDeals)
.where(eq(supplyDeals.id, record.supplyDealId))
.limit(1);
if (!deal) throw MutationError.notFound('Supply deal');
if (record.accountId && deal.accountId !== record.accountId) {
throw new MutationError(
'relationship_mismatch',
'Supply deal belongs to a different account.',
409,
);
}
}
}
export function createEntryMutationDefinition(): MutationDefinition<
typeof createEntrySchema,
CalendarEntry
> {
return {
schema: createEntrySchema,
permission: { authorize: requireCalendarWrite },
invalidMessage: 'Invalid calendar entry.',
async mutate({ input, principal, tx, now }) {
const values = {
...input,
startsAt: new Date(input.startsAt),
endsAt: date(input.endsAt) ?? null,
completedAt: date(input.completedAt) ?? null,
// Unassigned work is work nobody does, so an entry defaults to the
// person creating it rather than to nobody.
ownerUserId: input.ownerUserId === undefined ? principal.userId : input.ownerUserId,
createdByUserId: principal.userId,
updatedAt: now,
};
await requireRelationships(tx, { ...values, ownerUserId: input.ownerUserId ?? null });
const created = writtenRow(
(await tx.insert(calendarEntries).values(values).returning())[0],
);
return {
data: created,
activity: {
type: created.kind === 'meeting' || created.kind === 'qbr' ? 'meeting' : 'task',
subject: `Scheduled ${created.title}`,
accountId: created.accountId ?? undefined,
demandDealId: created.demandDealId ?? undefined,
supplyDealId: created.supplyDealId ?? undefined,
meta: {
calendarEntryId: created.id,
entryKind: created.kind,
startsAt: created.startsAt,
},
},
};
},
};
}
export function updateEntryMutationDefinition(): MutationDefinition<
typeof updateEntrySchema,
CalendarEntry
> {
return {
schema: updateEntrySchema,
permission: { authorize: requireCalendarWrite },
invalidMessage: 'Invalid calendar entry update.',
async mutate({ input, params, tx, now }) {
const id = requiredRouteParam(params);
const [before] = await tx
.select()
.from(calendarEntries)
.where(eq(calendarEntries.id, id))
.limit(1);
if (!before) throw MutationError.notFound('Calendar entry');
const changes = {
...input,
startsAt: input.startsAt ? new Date(input.startsAt) : undefined,
endsAt: date(input.endsAt),
completedAt: date(input.completedAt),
updatedAt: now,
};
// `undefined` means "leave alone" and `null` means "clear", so the row
// being validated has to be the merge, not the patch.
await requireRelationships(tx, {
// Unlike the others this is the patch, not the merge: the stored owner
// was checked when it was written and may since have been deleted, and
// failing an unrelated edit over that helps nobody.
ownerUserId: changes.ownerUserId ?? null,
accountId: changes.accountId === undefined ? before.accountId : changes.accountId,
demandDealId:
changes.demandDealId === undefined ? before.demandDealId : changes.demandDealId,
supplyDealId:
changes.supplyDealId === undefined ? before.supplyDealId : changes.supplyDealId,
startsAt: changes.startsAt ?? before.startsAt,
endsAt: changes.endsAt === undefined ? before.endsAt : changes.endsAt,
});
const updated = writtenRow(
(
await tx
.update(calendarEntries)
.set(changes)
.where(eq(calendarEntries.id, before.id))
.returning()
)[0],
);
return {
data: updated,
activity: {
type: 'task',
subject: `${updated.completedAt ? 'Completed' : 'Updated'} ${updated.title}`,
accountId: updated.accountId ?? undefined,
demandDealId: updated.demandDealId ?? undefined,
supplyDealId: updated.supplyDealId ?? undefined,
meta: {
calendarEntryId: updated.id,
completedAt: updated.completedAt,
startsAt: updated.startsAt,
},
},
};
},
};
}
export function deleteEntryMutationDefinition(): MutationDefinition<
z.ZodObject<Record<string, never>>,
{ id: string }
> {
return {
schema: z.object({}),
permission: { authorize: requireCalendarWrite },
invalidMessage: 'Invalid calendar entry deletion.',
async mutate({ params, tx }) {
const id = requiredRouteParam(params);
const [deleted] = await tx
.delete(calendarEntries)
.where(eq(calendarEntries.id, id))
.returning();
if (!deleted) throw MutationError.notFound('Calendar entry');
return {
data: { id: deleted.id },
activity: {
type: 'task',
subject: `Removed ${deleted.title}`,
accountId: deleted.accountId ?? undefined,
demandDealId: deleted.demandDealId ?? undefined,
supplyDealId: deleted.supplyDealId ?? undefined,
meta: { calendarEntryId: deleted.id, entryKind: deleted.kind },
},
};
},
};
}
export const querySchema = z.object({
from: isoDate.optional(),
to: isoDate.optional(),
quarter: z.string().optional(),
kinds: z.string().optional(),
accountId: z.string().uuid().optional(),
ownerUserId: z.string().uuid().optional(),
/**
* Zero-based month index the fiscal year starts on. Passed per request
* because PIG has nowhere to store an organisation-wide fiscal calendar yet,
* and inventing a settings column here would be a second place for the
* answer to live. Calendar quarters remain the default.
*/
fiscalYearStartMonth: z.coerce.number().int().min(0).max(11).optional(),
/**
* Checked against ICU here rather than left to the UTC fallback in
* `@pig/core`: the fallback exists for `users.timezone`, which is already
* stored and cannot be argued with, whereas a caller who asked for
* `Mars/Olympus` can be told. It also keeps the formatter cache — keyed on
* this string — from being fed arbitrary values by a caller in a loop.
*/
timezone: z
.string()
.max(80)
.refine(isValidTimeZone, { message: 'Unknown IANA time zone.' })
.optional(),
});
/** The one filter the entries listing takes; a uuid column cannot be asked about free text. */
export const entriesQuerySchema = z.object({
accountId: z.string().uuid().optional(),
});
export function createCalendarRoutes(db: Database): Hono<ApiEnv> {
const routes = new Hono<ApiEnv>();
const service = new CalendarService(db);
routes.get('/api/calendar', async (context: Context<ApiEnv>) => {
const principal = context.get('principal');
if (!calendarReadAllowed(principal.scopes)) {
return context.json(
apiError('insufficient_scope', "This credential lacks the 'read' scope."),
403,
);
}
const parsed = querySchema.safeParse(context.req.query());
if (!parsed.success) {
return context.json(
apiError('invalid_query', 'Invalid calendar query.', parsed.error.issues),
400,
);
}
const query = parsed.data;
const fiscalYearStartMonth = query.fiscalYearStartMonth ?? 0;
// The reader's own zone decides where a quarter begins; an explicit
// parameter wins so a shared link shows both people the same window.
const timeZone = query.timezone ?? (await service.timeZoneFor(principal.userId));
let from: Date;
let to: Date;
if (query.from && query.to) {
from = new Date(query.from);
to = new Date(query.to);
if (to <= from) {
return context.json(apiError('invalid_range', "'to' must be after 'from'."), 400);
}
} else if (query.quarter) {
const label = parseQuarter(query.quarter);
if (!label) {
return context.json(
apiError('invalid_quarter', "Expected a quarter label such as '2026-Q3'."),
400,
);
}
const bounds = quarterBounds(label.year, label.quarter, fiscalYearStartMonth, timeZone);
from = bounds.from;
to = bounds.to;
} else if (query.from || query.to) {
return context.json(
apiError('invalid_range', "Provide both 'from' and 'to', or neither."),
400,
);
} else {
// No range at all is the common case — a GTM lead opening the page wants
// the quarter they are standing in.
const bounds = quarterBoundsFor(new Date(), fiscalYearStartMonth, timeZone);
from = bounds.from;
to = bounds.to;
}
let kinds: CalendarEventKind[] | undefined;
try {
kinds = parseKinds(query.kinds);
} catch (error) {
if (error instanceof MutationError) {
return context.json(apiError(error.code, error.message), error.status);
}
throw error;
}
return context.json(
await service.project({
from,
to,
kinds,
accountId: query.accountId,
ownerUserId: query.ownerUserId,
fiscalYearStartMonth,
timeZone,
}),
);
});
/** The owned rows, listed on their own so the CRUD is inspectable. */
routes.get('/api/calendar/entries', async (context: Context<ApiEnv>) => {
const principal = context.get('principal');
if (!calendarReadAllowed(principal.scopes)) {
return context.json(
apiError('insufficient_scope', "This credential lacks the 'read' scope."),
403,
);
}
// Postgres rejects a malformed uuid with 22P02, which surfaces as a 500;
// the sibling read above already answers 400 for the same parameter.
const parsed = entriesQuerySchema.safeParse(context.req.query());
if (!parsed.success) {
return context.json(
apiError('invalid_query', 'Invalid calendar entries query.', parsed.error.issues),
400,
);
}
const accountId = parsed.data.accountId;
const rows = await db
.select({ entry: calendarEntries, accountName: accounts.name })
.from(calendarEntries)
.leftJoin(accounts, eq(accounts.id, calendarEntries.accountId))
.where(and(accountId ? eq(calendarEntries.accountId, accountId) : undefined))
.orderBy(calendarEntries.startsAt)
.limit(500);
return context.json({ kinds: CALENDAR_ENTRY_KINDS, entries: rows });
});
routes.post('/api/calendar/entries', mutation(db, createEntryMutationDefinition()));
routes.patch('/api/calendar/entries/:id', mutation(db, updateEntryMutationDefinition()));
routes.delete(
'/api/calendar/entries/:id',
bodylessMutation(db, deleteEntryMutationDefinition()),
);
return routes;
}
+5 -1
View File
@@ -45,7 +45,11 @@ export const factDecisionDefinition: MutationDefinition<
FactDecisionResult
> = {
schema: factDecisionSchema,
permission: { capability: 'data:import', team: 'research' },
// `fact:review`, not `data:import`. Accepting an agent's claim about a named
// person is a judgement about evidence; rewriting five thousand rows from a
// spreadsheet is not. They shared a capability until an audit noticed that
// granting either granted both.
permission: { capability: 'fact:review', team: 'research' },
invalidMessage: 'Invalid fact review decision.',
async mutate({ input, params, principal, tx, now }) {
const id = params.id;
+23 -2
View File
@@ -1,7 +1,7 @@
import type { Database } from '@pig/db';
import { Hono } from 'hono';
import { z } from 'zod';
import { requireCapability } from '../lib/auth';
import { requireAnyTeamCapability } from '../lib/auth';
import type { ApiEnv } from '../lib/mutation';
import { MutationError } from '../lib/mutation';
import {
@@ -50,8 +50,29 @@ export function createGoogleSheetsRoutes(
}
});
/*
* Two capabilities, not one.
*
* Handing PIG a long-lived Google refresh token is `integration:connect`:
* an authority over a third-party account, granted once, revocable
* separately. Reading the resulting spreadsheets in order to import them is
* `data:import`. Someone allowed to connect their Drive is not thereby
* allowed to rewrite the book from it, and the reverse is just as true.
*
* `data:import` is any-team here because no import entity has been chosen
* yet — the browser is still picking a file. The team is enforced at commit,
* in imports.ts, where the target is known.
*/
routes.use('/api/imports/google/*', async (context, next) => {
requireCapability(context.get('principal'), 'data:import');
const path = new URL(context.req.url).pathname;
const managesConnection =
path === '/api/imports/google/status' ||
path === '/api/imports/google/connect' ||
path === '/api/imports/google/connection';
requireAnyTeamCapability(
context.get('principal'),
managesConnection ? 'integration:connect' : 'data:import',
);
await next();
});
routes.get('/api/imports/google/status', async (context) =>
+10 -19
View File
@@ -1,5 +1,5 @@
import type { Database } from '@pig/db';
import { Hono, type Context } from 'hono';
import { Hono } from 'hono';
import { z } from 'zod';
import type { ApiEnv } from '../lib/mutation';
import { apiError } from '../lib/mutation';
@@ -7,29 +7,20 @@ import { CustomerLifecycleService } from '../services/customer-lifecycle';
const accountIdSchema = z.string().uuid();
export function growthReadAllowed(scopes: readonly string[]): boolean {
return scopes.includes('read');
}
function requireGrowthRead(context: Context<ApiEnv>) {
if (growthReadAllowed(context.get('principal').scopes)) return null;
return context.json(
apiError('insufficient_scope', "This credential lacks the 'read' scope."),
403,
);
}
/*
* These used to carry their own `requireGrowthRead`, which asked whether the
* *credential* had the 'read' scope. That was the right instinct and the wrong
* question: scope is a property of the API key, and it said nothing about
* whether the person holding it may see the growth book. Both halves are now
* asked once, for every read in the product, by the READ_RULES table —
* `requireReadCapability` checks the scope first and then `book:read`.
*/
export function createGrowthRoutes(db: Database): Hono<ApiEnv> {
const routes = new Hono<ApiEnv>();
const service = new CustomerLifecycleService(db);
routes.get('/api/growth', async (context) => {
const denied = requireGrowthRead(context);
return denied ?? context.json(await service.report());
});
routes.get('/api/growth', async (context) => context.json(await service.report()));
routes.get('/api/growth/accounts/:id', async (context) => {
const denied = requireGrowthRead(context);
if (denied) return denied;
const accountId = accountIdSchema.safeParse(context.req.param('id'));
if (!accountId.success) {
return context.json(apiError('invalid_account', 'Invalid account ID.', accountId.error.issues), 400);
+5 -2
View File
@@ -3,7 +3,7 @@ import { HUBSPOT_OBJECT_TYPES } from '../integrations/hubspot/contracts';
import { HubSpotOAuthError } from '../integrations/hubspot/oauth';
import { Hono } from 'hono';
import { z } from 'zod';
import { requireCapability } from '../lib/auth';
import { requireAnyTeamCapability, requireCapability } from '../lib/auth';
import type { ApiEnv } from '../lib/mutation';
const connectionParamSchema = z.string().uuid();
@@ -70,7 +70,10 @@ export function createHubSpotRoutes(service: HubSpotRouteService): Hono<ApiEnv>
routes.post('/api/integrations/hubspot/connections/:connectionId/sync', async (context) => {
const principal = context.get('principal');
requireCapability(principal, 'data:import');
// Any team, and honestly so: a HubSpot sync pulls companies, contacts and
// deals from both sides at once, so there is no single team to scope it to.
// Narrowing it would need the sync to accept an object-type filter first.
requireAnyTeamCapability(principal, 'data:import');
const parsedId = connectionParamSchema.safeParse(context.req.param('connectionId'));
if (!parsedId.success) return context.json({ error: 'Invalid HubSpot connection ID.' }, 400);
return context.json(
+54 -4
View File
@@ -1,8 +1,8 @@
import { IMPORT_ENTITIES, IMPORT_ENTITY_DEFINITIONS } from '@pig/core';
import { IMPORT_ENTITIES, IMPORT_ENTITY_DEFINITIONS, type ImportEntity, type Team } from '@pig/core';
import type { Database } from '@pig/db';
import { Hono } from 'hono';
import { z } from 'zod';
import { requireCapability } from '../lib/auth';
import { requireAnyTeamCapability, requireCapability } from '../lib/auth';
import type { ApiEnv, MutationDefinition } from '../lib/mutation';
import { MutationError, mutation } from '../lib/mutation';
import {
@@ -37,6 +37,39 @@ const parseSchema = z.object({
base64: z.string().min(1).max(Math.ceil(MAX_IMPORT_FILE_BYTES * 4 / 3) + 16),
}).strict();
/**
* Which team's book an import writes into.
*
* The commit is the only point at which that is knowable, and it is the only
* point at which it matters: `requireCapability(p, 'data:import')` with no team
* passed if the principal held the capability on *any* team, so a research-team
* admin could rewrite the demand pipeline. Accounts and contacts are shared by
* both commercial sides, so admin of either is enough for those.
*/
export const IMPORT_ENTITY_TEAMS: Readonly<Record<ImportEntity, readonly Team[]>> = {
account: ['supply', 'demand'],
contact: ['supply', 'demand'],
demand_deal: ['demand'],
supply_deal: ['supply'],
};
function requireImportPermission(
principal: Parameters<typeof requireCapability>[0],
entity: ImportEntity,
): void {
const teams = IMPORT_ENTITY_TEAMS[entity];
let denial: unknown;
for (const team of teams) {
try {
requireCapability(principal, 'data:import', team);
return;
} catch (error) {
denial = error;
}
}
throw denial;
}
interface ImportCommitOperations {
commit(
input: z.infer<typeof commitSchema>,
@@ -52,9 +85,13 @@ export function createImportCommitMutationDefinition(
): MutationDefinition<typeof commitSchema, ImportCommitResult> {
return {
schema: commitSchema,
permission: { authorize: (principal) => requireCapability(principal, 'data:import') },
// Two stages: any-team up front so a principal with no import authority at
// all cannot probe the schema, then the entity's own team once the body has
// been parsed and the target is finally knowable.
permission: { authorize: (principal) => requireAnyTeamCapability(principal, 'data:import') },
invalidMessage: 'Invalid import commit.',
async mutate({ input, principal, tx, now }) {
requireImportPermission(principal, input.entity);
const result = await makeService(tx).commit(input, principal, now);
const entityLabel = IMPORT_ENTITY_DEFINITIONS[input.entity].label.toLocaleLowerCase();
return {
@@ -78,8 +115,21 @@ export function createImportCommitMutationDefinition(
export function createImportRoutes(db: Database): Hono<ApiEnv> {
const routes = new Hono<ApiEnv>();
// Any team, deliberately: config, parse and preview touch no book at all —
// preview is a dry run against uploaded cells. The commit is where the team
// is enforced, because the commit is where rows are written.
//
// The nested integration namespaces are skipped rather than left to fall
// through this. They are mounted after this router, so Hono runs this
// middleware for them too, and it would have re-imposed `data:import` on the
// OAuth routes that were just split onto `integration:connect` — the split
// would have compiled, passed its unit tests, and changed nothing.
routes.use('/api/imports/*', async (context, next) => {
requireCapability(context.get('principal'), 'data:import');
const path = new URL(context.req.url).pathname;
if (path.startsWith('/api/imports/google/') || path.startsWith('/api/imports/notion/')) {
return next();
}
requireAnyTeamCapability(context.get('principal'), 'data:import');
await next();
});
routes.get('/api/imports/config', (context) => context.json({
+782
View File
@@ -0,0 +1,782 @@
/**
* Learn — the member curriculum, the admin CRUD, and the one door in this API
* that opens without a principal.
*
* ## The security shape, which is the reason this file is long
*
* Everywhere else in PIG a request resolves a `Principal` and then a
* capability check decides what it may do. A code-holder has no account, so
* there is no principal to resolve — and the tempting shortcut, minting a
* synthetic one, is the thing this design exists to refuse. A principal is
* accepted by every downstream handler by construction; the only thing keeping
* it out of the CRM would be that each of those handlers remembered to check a
* capability. One that forgot would leak the book of business to anyone
* holding a marketing share code, and nothing would report an error.
*
* So the code mints a **scoped bearer token that is not a credential for this
* API at all**. It is an HMAC over a scope string and an expiry, verified by
* exactly one handler, and `authenticate()` never sees it. Presenting it to
* `/api/dashboard` produces the same 401 as presenting nothing, because to the
* authenticator it is simply a bearer token that is not a JWT and does not
* start with `pig_`. There is a test that asserts precisely this.
*
* The token's signing key is derived from the stored access code, so rotating
* the code invalidates every outstanding token as a side effect rather than
* requiring a second revocation mechanism.
*
* ## The public read
*
* `GET /api/learn/public` is written so that it is *structurally* incapable of
* naming another row: the predicates are two literals, there is no parameter
* that reaches the WHERE clause, and the selected columns are enumerated. It
* cannot be widened by a query string because it does not read one.
*
* ## What may be code-visible
*
* Only the platform track. Enforced here in the write path AND by a CHECK
* constraint on the table. Not in the UI: a form is not a security boundary,
* and a concept video becoming anon-visible through a mis-set select is the
* failure that matters.
*/
import { and, asc, desc, eq, isNull } from 'drizzle-orm';
import { Hono } from 'hono';
import { createHash, createHmac, timingSafeEqual } from 'node:crypto';
import { z } from 'zod';
import {
LEARN_CODE_TRACK,
LEARN_EMBED_REJECTION_MESSAGES,
LEARN_TRACKS,
LEARN_VISIBILITIES,
learnEmbedUrl,
learnVisibilityPermitted,
learnWatchUrl,
resolveLearnEmbed,
type LearnProvider,
type LearnTrack,
type LearnVisibility,
} from '@pig/core';
import { learnResources, platformSettings } from '@pig/db';
import type { Database } from '@pig/db';
import { requireCapability, safeEqual } from '../lib/auth';
import {
apiError,
bodylessMutation,
MutationError,
mutation,
type ApiEnv,
} from '../lib/mutation';
/**
* The two paths that must be allowlisted in `app.ts`, exported as constants
* for the same reason the Slack and Notion callbacks are: a public path
* spelled twice is a public path that eventually differs in one of them.
*/
export const LEARN_ACCESS_PATH = '/api/learn/access';
export const LEARN_PUBLIC_PATH = '/api/learn/public';
const SETTINGS_ID = 'default';
// ---------------------------------------------------------------------- token
const LEARN_TOKEN_VERSION = 'v1';
/**
* Baked into the signature, not merely into the format. A token is a claim
* about a scope; if the scope were only in the envelope, widening the format
* later would silently promote every token already in a browser.
*/
const LEARN_TOKEN_SCOPE = `learn:${LEARN_CODE_TRACK}`;
export const LEARN_TOKEN_PREFIX = 'learn_';
/**
* Long enough that someone working through onboarding is not interrupted,
* short enough that a code rotation is not the only way to end a session. The
* token grants nothing but the platform track, so the usual argument for a
* short expiry — blast radius — barely applies.
*/
export const LEARN_TOKEN_TTL_MS = 12 * 60 * 60 * 1_000;
/**
* The signing key, derived from the code rather than configured separately.
*
* This is what makes rotation total: change the code and every token already
* in a browser stops verifying, with no revocation list to maintain and no
* second secret to keep in step. The separator keeps the salt and the code
* from running together, so no two codes can yield the same key material.
*/
function tokenKey(accessCode: string): Buffer {
return createHash('sha256')
.update(`pig.learn.token.${LEARN_TOKEN_VERSION}\u0000${accessCode}`)
.digest();
}
export function mintLearnToken(accessCode: string, expiresAt: number): string {
const signature = createHmac('sha256', tokenKey(accessCode))
.update(`${LEARN_TOKEN_VERSION}.${LEARN_TOKEN_SCOPE}.${expiresAt}`)
.digest('base64url');
return `${LEARN_TOKEN_PREFIX}${LEARN_TOKEN_VERSION}.${expiresAt}.${signature}`;
}
export const LEARN_TOKEN_REJECTIONS = ['missing', 'malformed', 'expired', 'mismatch'] as const;
export type LearnTokenRejection = (typeof LEARN_TOKEN_REJECTIONS)[number];
export type LearnTokenResult =
| { valid: true; expiresAt: number }
| { valid: false; reason: LearnTokenRejection };
/**
* Verify a learn token against the code currently in force.
*
* Follows the register of `integrations/hubspot/signature.ts`: recompute,
* compare byte lengths first because `timingSafeEqual` throws on a mismatch,
* then compare in constant time. Expiry is checked before the HMAC only
* because an expired token is not a secret worth protecting the timing of.
*/
export function verifyLearnToken(
accessCode: string | null,
token: string | undefined,
now: number = Date.now(),
): LearnTokenResult {
if (!accessCode) return { valid: false, reason: 'mismatch' };
if (!token) return { valid: false, reason: 'missing' };
if (!token.startsWith(LEARN_TOKEN_PREFIX)) return { valid: false, reason: 'malformed' };
const parts = token.slice(LEARN_TOKEN_PREFIX.length).split('.');
if (parts.length !== 3) return { valid: false, reason: 'malformed' };
const [version, rawExpiry, signature] = parts as [string, string, string];
if (version !== LEARN_TOKEN_VERSION) return { valid: false, reason: 'malformed' };
if (!/^\d{1,15}$/.test(rawExpiry)) return { valid: false, reason: 'malformed' };
const expiresAt = Number(rawExpiry);
if (!Number.isSafeInteger(expiresAt)) return { valid: false, reason: 'malformed' };
if (expiresAt <= now) return { valid: false, reason: 'expired' };
const expected = Buffer.from(
createHmac('sha256', tokenKey(accessCode))
.update(`${LEARN_TOKEN_VERSION}.${LEARN_TOKEN_SCOPE}.${expiresAt}`)
.digest('base64url'),
'utf8',
);
const supplied = Buffer.from(signature, 'utf8');
if (expected.length !== supplied.length) return { valid: false, reason: 'mismatch' };
return timingSafeEqual(expected, supplied)
? { valid: true, expiresAt }
: { valid: false, reason: 'mismatch' };
}
// -------------------------------------------------------------- rate limiting
export interface AttemptDecision {
allowed: boolean;
remaining: number;
retryAfterSeconds: number;
}
export interface AttemptLimiter {
check(key: string, now?: number): AttemptDecision;
}
/**
* A fixed-window limiter, in process.
*
* Deliberately not distributed and deliberately not durable. PIG runs as one
* container; the job here is to stop a script walking a short passphrase
* keyspace at HTTP speed, not to enforce an exact quota. A restart resetting
* the window costs an attacker one restart's worth of guesses, which is not
* the difference between safe and unsafe — the code length is.
*
* The map is pruned on write rather than on a timer, and capped, because the
* key is a client-supplied-ish address and an unbounded map keyed on one is a
* memory-exhaustion primitive.
*/
export function createAttemptLimiter({
limit,
windowMs,
maxKeys = 10_000,
}: {
limit: number;
windowMs: number;
maxKeys?: number;
}): AttemptLimiter {
const windows = new Map<string, { count: number; resetAt: number }>();
return {
check(key, now = Date.now()) {
if (windows.size >= maxKeys) {
for (const [existing, window] of windows) {
if (window.resetAt <= now) windows.delete(existing);
}
// Still full: every window is live, so this is either a real flood or
// a spoofed-address one. Refuse rather than grow.
if (windows.size >= maxKeys) {
return { allowed: false, remaining: 0, retryAfterSeconds: Math.ceil(windowMs / 1000) };
}
}
const current = windows.get(key);
if (!current || current.resetAt <= now) {
windows.set(key, { count: 1, resetAt: now + windowMs });
return { allowed: true, remaining: limit - 1, retryAfterSeconds: 0 };
}
current.count += 1;
if (current.count > limit) {
return {
allowed: false,
remaining: 0,
retryAfterSeconds: Math.max(1, Math.ceil((current.resetAt - now) / 1000)),
};
}
return { allowed: true, remaining: limit - current.count, retryAfterSeconds: 0 };
},
};
}
/**
* Which client is this, for rate-limiting purposes?
*
* The LAST entry in `X-Forwarded-For`, not the first. Caddy APPENDS the real
* peer to whatever the client sent, so the first hop is attacker-controlled
* and using it hands anyone an unlimited number of rate-limit buckets. Behind
* exactly one proxy — which is this deployment — the last entry is the only
* one the client could not write.
*/
export function rateLimitKey(forwardedFor: string | undefined): string {
if (!forwardedFor) return 'unknown';
const hops = forwardedFor
.split(',')
.map((hop) => hop.trim())
.filter(Boolean);
return hops[hops.length - 1] ?? 'unknown';
}
// ------------------------------------------------------------------- schemas
const urlField = z.string().trim().min(1).max(2_000);
export const learnAccessSchema = z
.object({ code: z.string().min(1).max(200) })
.strict();
export const learnResourceCreateSchema = z
.object({
track: z.enum(LEARN_TRACKS),
title: z.string().trim().min(1).max(200),
summary: z.string().trim().max(1_000).optional(),
url: urlField,
visibility: z.enum(LEARN_VISIBILITIES).default('members'),
// A day is generous for a walkthrough and rules out a mistyped
// milliseconds value being stored as seconds.
durationSeconds: z.number().int().positive().max(86_400).optional(),
sortOrder: z.number().int().min(0).max(10_000).optional(),
publishedAt: z.string().datetime().optional(),
})
.strict()
.refine(
(value) => learnVisibilityPermitted(value.track, value.visibility),
'Only platform-track resources may be unlocked by the share code.',
);
export const learnResourceUpdateSchema = z
.object({
track: z.enum(LEARN_TRACKS).optional(),
title: z.string().trim().min(1).max(200).optional(),
summary: z.string().trim().max(1_000).nullable().optional(),
url: urlField.optional(),
visibility: z.enum(LEARN_VISIBILITIES).optional(),
durationSeconds: z.number().int().positive().max(86_400).nullable().optional(),
sortOrder: z.number().int().min(0).max(10_000).optional(),
publishedAt: z.string().datetime().optional(),
archived: z.boolean().optional(),
})
.strict()
.refine((value) => Object.values(value).some((item) => item !== undefined), 'No changes supplied.');
/**
* A passphrase a human reads aloud, so printable ASCII and no whitespace.
* Six is the floor because the endpoint is rate-limited, not because a short
* code is otherwise fine.
*/
export const learnAccessCodeSchema = z
.object({ code: z.string().trim().min(6).max(120).regex(/^[\x21-\x7e]+$/, 'Use printable characters with no spaces.') })
.strict();
// ------------------------------------------------------------- serialisation
interface LearnRowForView {
id: string;
track: LearnTrack;
title: string;
summary: string | null;
provider: LearnProvider;
externalId: string;
visibility: LearnVisibility;
durationSeconds: number | null;
sortOrder: number;
publishedAt: Date;
}
/**
* The wire shape. Note what is absent: the stored `url` column never leaves
* the database. Both URLs a client receives are rebuilt from the allowlist,
* so a row whose `url` was poisoned by some future write path still cannot put
* an attacker's bytes into an `iframe src`.
*
* A row we cannot rebuild an embed for is dropped rather than returned
* without one — it would render as a card that does nothing, and the honest
* reading of "this provider is no longer enabled" is that the video is not
* available, not that it is broken.
*/
export function learnResourceView(row: LearnRowForView) {
const embedUrl = learnEmbedUrl(row.provider, row.externalId);
const watchUrl = learnWatchUrl(row.provider, row.externalId);
if (!embedUrl || !watchUrl) return null;
return {
id: row.id,
track: row.track,
title: row.title,
summary: row.summary,
provider: row.provider,
visibility: row.visibility,
durationSeconds: row.durationSeconds,
sortOrder: row.sortOrder,
publishedAt: row.publishedAt,
embedUrl,
watchUrl,
};
}
export type LearnResourceView = NonNullable<ReturnType<typeof learnResourceView>>;
function renderable(rows: LearnRowForView[]): LearnResourceView[] {
return rows.map(learnResourceView).filter((view): view is LearnResourceView => view !== null);
}
/** Enumerated rather than `select()`, so `url` cannot be added by accident. */
const viewColumns = {
id: learnResources.id,
track: learnResources.track,
title: learnResources.title,
summary: learnResources.summary,
provider: learnResources.provider,
externalId: learnResources.externalId,
visibility: learnResources.visibility,
durationSeconds: learnResources.durationSeconds,
sortOrder: learnResources.sortOrder,
publishedAt: learnResources.publishedAt,
} as const;
// ---------------------------------------------------------------------- routes
async function currentAccessCode(db: Database): Promise<string | null> {
const [row] = await db
.select({ code: platformSettings.learnAccessCode })
.from(platformSettings)
.where(eq(platformSettings.id, SETTINGS_ID))
.limit(1);
// No settings row means the workspace has not been initialised. Refusing
// every code is the correct answer; inserting a row here would write default
// Piggy configuration over what `ensurePlatformSettings` derives from the
// environment, which is a far worse bug than an unusable share code.
return row?.code ?? null;
}
export function createLearnRoutes(
db: Database,
options: { limiter?: AttemptLimiter } = {},
) {
const app = new Hono<ApiEnv>();
// Ten guesses a minute per address. A human who has been given the code
// types it once; anything approaching this rate is a script.
const limiter = options.limiter ?? createAttemptLimiter({ limit: 10, windowMs: 60_000 });
// ------------------------------------------------------------- public door
app.post(LEARN_ACCESS_PATH, async (c) => {
const decision = limiter.check(rateLimitKey(c.req.header('x-forwarded-for')));
if (!decision.allowed) {
c.header('retry-after', String(decision.retryAfterSeconds));
return c.json(
apiError('learn_rate_limited', 'Too many attempts. Try again shortly.'),
429,
);
}
let body: unknown;
try {
body = await c.req.json();
} catch {
return c.json(apiError('invalid_json', 'Request body must be valid JSON.'), 400);
}
const parsed = learnAccessSchema.safeParse(body);
if (!parsed.success) {
return c.json(apiError('invalid_request', 'Supply an access code.', parsed.error.issues), 400);
}
const accessCode = await currentAccessCode(db);
if (!accessCode || !safeEqual(parsed.data.code.trim(), accessCode)) {
// One message for "wrong code" and "no code configured". Distinguishing
// them tells a guesser whether to keep going.
return c.json(apiError('invalid_code', 'That code is not valid.'), 401);
}
const expiresAt = Date.now() + LEARN_TOKEN_TTL_MS;
return c.json({
token: mintLearnToken(accessCode, expiresAt),
expiresAt: new Date(expiresAt).toISOString(),
/**
* Stated in the response because the front end has to be able to explain
* to a code-holder why the Concepts section is locked, and hardcoding
* that in the browser would be a second place to change it.
*/
track: LEARN_CODE_TRACK,
});
});
/**
* The only read a non-member can perform.
*
* Two literal predicates and no parameters. There is nothing in this handler
* that a caller can influence except the token, which decides whether it
* runs at all — not what it returns.
*/
app.get(LEARN_PUBLIC_PATH, async (c) => {
const header = c.req.header('authorization');
const supplied = header?.startsWith('Bearer ') ? header.slice(7).trim() : undefined;
const result = verifyLearnToken(await currentAccessCode(db), supplied);
if (!result.valid) {
return c.json(
apiError(
result.reason === 'expired' ? 'learn_token_expired' : 'learn_token_invalid',
result.reason === 'expired'
? 'That access has expired. Enter the code again.'
: 'A valid access code is required.',
),
401,
);
}
const rows = await db
.select(viewColumns)
.from(learnResources)
.where(
and(
eq(learnResources.visibility, 'code'),
eq(learnResources.track, LEARN_CODE_TRACK),
isNull(learnResources.archivedAt),
),
)
.orderBy(asc(learnResources.sortOrder), desc(learnResources.publishedAt));
return c.json({
track: LEARN_CODE_TRACK,
expiresAt: new Date(result.expiresAt).toISOString(),
resources: renderable(rows),
/** So the locked Concepts panel can name what is behind it. */
lockedTracks: LEARN_TRACKS.filter((track) => track !== LEARN_CODE_TRACK),
});
});
// ------------------------------------------------------------ member reads
/**
* The whole curriculum. Any member may read it: this is training material,
* and gating supply concepts behind supply-team membership would stop a new
* demand seller learning how the other side works, which is the opposite of
* what the page is for.
*/
app.get('/api/learn', async (c) => {
const rows = await db
.select(viewColumns)
.from(learnResources)
.where(isNull(learnResources.archivedAt))
.orderBy(asc(learnResources.sortOrder), desc(learnResources.publishedAt));
const views = renderable(rows);
return c.json({
tracks: Object.fromEntries(
LEARN_TRACKS.map((track) => [track, views.filter((view) => view.track === track)]),
) as Record<LearnTrack, LearnResourceView[]>,
canManage: canManageLearn(c.get('principal')),
});
});
// ------------------------------------------------------------- admin write
app.post(
'/api/learn/resources',
mutation(db, {
schema: learnResourceCreateSchema,
permission: { capability: 'settings:admin' },
invalidMessage: 'Invalid learn resource.',
async mutate({ input, principal, tx, now }) {
const resolved = resolveEmbedOrThrow(input.url);
const [created] = await tx
.insert(learnResources)
.values({
track: input.track,
title: input.title,
summary: input.summary ?? null,
// The canonical form from the allowlist, not the pasted string —
// so the stored value is one we generated even in the column
// nothing renders.
url: resolved.watchUrl,
provider: resolved.provider,
externalId: resolved.externalId,
visibility: input.visibility,
durationSeconds: input.durationSeconds ?? null,
sortOrder: input.sortOrder ?? 100,
publishedAt: input.publishedAt ? new Date(input.publishedAt) : now,
addedByUserId: principal.userId,
createdAt: now,
updatedAt: now,
})
.returning();
if (!created) throw new Error('Learn resource insert returned no row');
return {
data: viewOrThrow(created),
activity: {
type: 'agent_action',
subject: `Added learn resource: ${created.title}`,
meta: {
action: 'learn_resource.created',
resourceId: created.id,
track: created.track,
visibility: created.visibility,
provider: created.provider,
},
},
};
},
}),
);
app.patch(
'/api/learn/resources/:id',
mutation(db, {
schema: learnResourceUpdateSchema,
permission: { capability: 'settings:admin' },
invalidMessage: 'Invalid learn resource change.',
async mutate({ input, params, tx, now }) {
const id = requiredId(params);
const [existing] = await tx
.select()
.from(learnResources)
.where(eq(learnResources.id, id))
.limit(1);
if (!existing) throw MutationError.notFound('Learn resource');
/*
* Checked against the MERGED row, not the input. A PATCH that sets
* only `visibility: 'code'` on a supply resource carries no track at
* all, so validating the input alone would wave it straight through
* into the CHECK constraint and a 500.
*/
const track = input.track ?? existing.track;
const visibility = input.visibility ?? existing.visibility;
if (!learnVisibilityPermitted(track, visibility)) {
throw new MutationError(
'visibility_not_permitted',
`Only ${LEARN_CODE_TRACK}-track resources may be unlocked by the share code.`,
409,
);
}
const set: Partial<typeof learnResources.$inferInsert> = { updatedAt: now };
if (input.track !== undefined) set.track = input.track;
if (input.title !== undefined) set.title = input.title;
if (input.summary !== undefined) set.summary = input.summary;
if (input.visibility !== undefined) set.visibility = input.visibility;
if (input.durationSeconds !== undefined) set.durationSeconds = input.durationSeconds;
if (input.sortOrder !== undefined) set.sortOrder = input.sortOrder;
if (input.publishedAt !== undefined) set.publishedAt = new Date(input.publishedAt);
if (input.archived !== undefined) set.archivedAt = input.archived ? now : null;
if (input.url !== undefined) {
const resolved = resolveEmbedOrThrow(input.url);
set.url = resolved.watchUrl;
set.provider = resolved.provider;
set.externalId = resolved.externalId;
}
const [updated] = await tx
.update(learnResources)
.set(set)
.where(eq(learnResources.id, id))
.returning();
if (!updated) throw MutationError.notFound('Learn resource');
return {
data: viewOrThrow(updated),
activity: {
type: 'agent_action',
subject: `Updated learn resource: ${updated.title}`,
meta: {
action: 'learn_resource.updated',
resourceId: updated.id,
fields: Object.keys(input),
track: updated.track,
visibility: updated.visibility,
},
},
};
},
}),
);
app.delete(
'/api/learn/resources/:id',
bodylessMutation(db, {
schema: z.object({}).strict(),
permission: { capability: 'settings:admin' },
invalidMessage: 'Invalid learn resource removal.',
// Archive, never delete: the activity log references the row, and "what
// did onboarding say in March?" is a real question.
async mutate({ params, tx, now }) {
const id = requiredId(params);
const [existing] = await tx
.select()
.from(learnResources)
.where(eq(learnResources.id, id))
.limit(1);
if (!existing) throw MutationError.notFound('Learn resource');
const [archived] = existing.archivedAt
? [existing]
: await tx
.update(learnResources)
.set({ archivedAt: now, updatedAt: now })
.where(eq(learnResources.id, id))
.returning();
if (!archived) throw MutationError.notFound('Learn resource');
return {
data: { id: archived.id, archivedAt: archived.archivedAt },
activity: {
type: 'agent_action',
subject: `Archived learn resource: ${archived.title}`,
meta: {
action: 'learn_resource.archived',
resourceId: archived.id,
alreadyArchived: existing.archivedAt !== null,
},
},
};
},
}),
);
// ------------------------------------------------------------ the code itself
/**
* Returned in clear to a platform administrator, deliberately.
*
* It is a passphrase they have to be able to read out to the person they are
* sharing a demo with — a code an admin cannot see is a code nobody can use.
* It is not a credential for anything but the platform track, and the read
* requires `settings:admin`.
*/
app.get('/api/learn/access-code', async (c) => {
requireCapability(c.get('principal'), 'settings:admin');
const [row] = await db
.select({
code: platformSettings.learnAccessCode,
updatedAt: platformSettings.learnAccessCodeUpdatedAt,
})
.from(platformSettings)
.where(eq(platformSettings.id, SETTINGS_ID))
.limit(1);
if (!row) {
return c.json(
apiError('settings_uninitialised', 'Open Settings once to initialise this workspace.'),
409,
);
}
return c.json({ code: row.code, updatedAt: row.updatedAt, url: '/learn' });
});
app.patch(
'/api/learn/access-code',
mutation(db, {
schema: learnAccessCodeSchema,
permission: { capability: 'settings:admin' },
invalidMessage: 'Invalid access code.',
async mutate({ input, tx, now }) {
/*
* UPDATE, never upsert. Inserting the row here would give it default
* Piggy configuration rather than the environment-derived values
* `ensurePlatformSettings` writes, silently disabling Piggy — a much
* worse outcome than telling an administrator to open Settings first.
*/
const [updated] = await tx
.update(platformSettings)
.set({ learnAccessCode: input.code, learnAccessCodeUpdatedAt: now, updatedAt: now })
.where(eq(platformSettings.id, SETTINGS_ID))
.returning({
code: platformSettings.learnAccessCode,
updatedAt: platformSettings.learnAccessCodeUpdatedAt,
});
if (!updated) {
throw new MutationError(
'settings_uninitialised',
'Open Settings once to initialise this workspace.',
409,
);
}
return {
data: updated,
activity: {
type: 'agent_action',
// The code itself is never written to the activity log: that log
// is readable by every member, and rotating a code into it would
// defeat the rotation.
subject: 'Rotated the Learn share code',
meta: { action: 'learn_access_code.rotated' },
},
};
},
}),
);
return app;
}
// ------------------------------------------------------------------- helpers
/** Curating the curriculum is an administrative act, not a GTM one. */
export function canManageLearn(principal: { isPlatformAdmin: boolean; scopes: string[] }): boolean {
return principal.isPlatformAdmin && principal.scopes.includes('write');
}
/**
* A row that has just passed `resolveEmbedOrThrow` must be renderable, so a
* null here is a contradiction between the resolver and the rebuilder rather
* than a resource that is merely unavailable. Fail loudly.
*/
function viewOrThrow(row: LearnRowForView): LearnResourceView {
const view = learnResourceView(row);
if (!view) throw new Error(`Learn resource ${row.id} was written but cannot be rendered`);
return view;
}
function requiredId(params: Readonly<Record<string, string>>): string {
const id = params.id;
if (!id) throw new MutationError('invalid_route_parameter', "Route parameter 'id' is required.", 400);
return id;
}
/**
* The write-path half of the allowlist. An unmatched URL is a 400, not a row —
* which is what makes "no unresolvable resource exists" an invariant rather
* than a hope.
*/
function resolveEmbedOrThrow(url: string) {
const resolved = resolveLearnEmbed(url);
if (!resolved.ok) {
throw new MutationError(
'invalid_video_url',
LEARN_EMBED_REJECTION_MESSAGES[resolved.reason],
400,
);
}
return resolved;
}
+13 -3
View File
@@ -4,7 +4,7 @@ import { deleteCookie, getCookie, setCookie } from 'hono/cookie';
import { Hono } from 'hono';
import { z } from 'zod';
import type { Config } from '../lib/config';
import { requireCapability } from '../lib/auth';
import { requireAnyTeamCapability } from '../lib/auth';
import type { ApiEnv } from '../lib/mutation';
import { decryptSecret, encryptSecret, encryptionReady } from '../lib/secrets';
import {
@@ -32,9 +32,19 @@ export function createNotionImportRoutes(
const oauthCookieName = config.isProduction ? '__Host-pig_notion_oauth' : 'pig_notion_oauth';
const oauthCookiePath = config.isProduction ? '/' : NOTION_OAUTH_CALLBACK_PATH;
/*
* Connecting a Notion workspace is `integration:connect`; materialising a
* data source into PIG rows is `data:import`. See the same split in
* google-sheets.ts for why they are not the same authority.
*/
routes.use('/api/imports/notion/*', async (context, next) => {
if (new URL(context.req.url).pathname === NOTION_OAUTH_CALLBACK_PATH) return next();
requireCapability(context.get('principal'), 'data:import');
const path = new URL(context.req.url).pathname;
if (path === NOTION_OAUTH_CALLBACK_PATH) return next();
const writesRows = path.endsWith('/materialize');
requireAnyTeamCapability(
context.get('principal'),
writesRows ? 'data:import' : 'integration:connect',
);
await next();
});
+59 -17
View File
@@ -1,7 +1,33 @@
import { PIGGY_PAGE_ROUTES, PIGGY_RECORD_TYPES } from '@pig/core';
import type { Database } from '@pig/db';
import { Hono } from 'hono';
import { stream } from 'hono/streaming';
import { z } from 'zod';
import type { Config } from '../lib/config';
import type { ApiEnv } from '../lib/mutation';
import { ensurePlatformSettings } from './admin-settings';
/**
* Derived from the @pig/core tuples, and kept in step with the identical
* schema in the Piggy chat server. Both are `.strict()`, so a context arm
* missing from either one is a 400 at that hop rather than a degraded answer.
*/
const contextSchema = z.discriminatedUnion('type', [
z
.object({
type: z.enum(PIGGY_RECORD_TYPES),
id: z.string().uuid(),
label: z.string().max(240).optional(),
})
.strict(),
z
.object({
type: z.literal('page'),
route: z.enum(PIGGY_PAGE_ROUTES),
label: z.string().max(240).optional(),
})
.strict(),
]);
const requestSchema = z
.object({
@@ -15,20 +41,7 @@ const requestSchema = z
)
.max(20)
.optional(),
context: z
.object({
type: z.enum([
'account',
'contact',
'demand_deal',
'supply_deal',
'contract',
'commitment',
]),
id: z.string().uuid(),
label: z.string().max(240).optional(),
})
.optional(),
context: contextSchema.optional(),
})
.strict();
@@ -37,15 +50,44 @@ export interface PiggyChatProxyOptions {
internalUrl?: string;
internalToken?: string;
fetchImpl?: typeof fetch;
/**
* The admin toggle, read per request. Omitted, the environment gate alone
* decides — which is what shipped, and why turning Piggy off in the admin UI
* did nothing.
*/
resolvePiggyEnabled?: () => Promise<boolean>;
}
/** The stored toggle. Paired with `createPiggyChatRoutes` at composition. */
export function platformPiggyEnabled(config: Config, db: Database): () => Promise<boolean> {
return async () => (await ensurePlatformSettings(config, db)).piggyEnabled;
}
export function createPiggyChatRoutes(options: PiggyChatProxyOptions) {
const routes = new Hono<ApiEnv>();
const fetchImpl = options.fetchImpl ?? fetch;
const available = Boolean(options.enabled && options.internalUrl && options.internalToken);
// Configuration cannot change under a running process; the toggle can.
const configured = Boolean(options.enabled && options.internalUrl && options.internalToken);
routes.get('/api/piggy/status', (c) => {
/**
* The environment variable is the outer gate and the stored setting the
* inner one: an operator who has not provisioned Piggy cannot have it
* switched on from the admin UI. A failed settings read falls back to the
* outer gate rather than 503-ing every dock on the site over one bad query.
*/
async function isAvailable(): Promise<boolean> {
if (!configured) return false;
if (!options.resolvePiggyEnabled) return true;
try {
return await options.resolvePiggyEnabled();
} catch {
return true;
}
}
routes.get('/api/piggy/status', async (c) => {
const principal = c.get('principal');
const available = await isAvailable();
return c.json({
enabled: available,
canUse: available && principal.scopes.includes('read'),
@@ -60,7 +102,7 @@ export function createPiggyChatRoutes(options: PiggyChatProxyOptions) {
403,
);
}
if (!available || !options.internalUrl || !options.internalToken) {
if (!(await isAvailable()) || !options.internalUrl || !options.internalToken) {
return c.json({ error: 'Piggy chat is not available.', code: 'piggy_unavailable' }, 503);
}
+60
View File
@@ -0,0 +1,60 @@
/**
* Which capability each read requires — the whole policy, in one table.
*
* It lives in a table rather than beside each handler because the question a
* reviewer needs to answer is "who can see cost?", and that question is
* unanswerable if the answer is spread across nine route files. Adding a GET
* without adding a row here leaves it ungoverned, which is the failure this
* exists to end; `read-governance.test.ts` fails when a new read path appears
* that no row covers.
*
* Mounted before every other route in `createApp`, and the order is
* load-bearing: Hono runs matched handlers in registration order, so a guard
* registered after its handler never runs.
*/
import type { ReadCapability } from '@pig/core';
import { Hono } from 'hono';
import { readGuard } from '../lib/read-guard';
import type { ApiEnv } from '../lib/mutation';
export interface ReadRule {
method: 'GET' | 'POST';
path: string;
capability: ReadCapability;
}
/**
* `economics:read` covers anything carrying supplier cost, break-even price or
* a margin total. `/api/capacity/match` is a POST only because a requirement
* is too big for a query string — it returns break-even per block, so it is a
* read and is gated as one.
*/
export const READ_RULES: readonly ReadRule[] = [
{ method: 'GET', path: '/api/capacity/availability', capability: 'economics:read' },
{ method: 'GET', path: '/api/capacity/idle', capability: 'economics:read' },
{ method: 'GET', path: '/api/capacity/margin', capability: 'economics:read' },
{ method: 'POST', path: '/api/capacity/match', capability: 'economics:read' },
{ method: 'GET', path: '/api/inventory', capability: 'economics:read' },
{ method: 'GET', path: '/api/commitments', capability: 'economics:read' },
{ method: 'GET', path: '/api/allocations', capability: 'economics:read' },
{ method: 'GET', path: '/api/dashboard', capability: 'economics:read' },
{ method: 'GET', path: '/api/accounts', capability: 'book:read' },
{ method: 'GET', path: '/api/accounts/:id', capability: 'book:read' },
{ method: 'GET', path: '/api/contacts', capability: 'book:read' },
{ method: 'GET', path: '/api/deals/demand', capability: 'book:read' },
{ method: 'GET', path: '/api/deals/supply', capability: 'book:read' },
{ method: 'GET', path: '/api/contracts', capability: 'book:read' },
{ method: 'GET', path: '/api/contracts/:id', capability: 'book:read' },
{ method: 'GET', path: '/api/growth', capability: 'book:read' },
{ method: 'GET', path: '/api/growth/accounts/:id', capability: 'book:read' },
{ method: 'GET', path: '/api/facts', capability: 'book:read' },
{ method: 'GET', path: '/api/team', capability: 'team:read' },
];
export function createReadGuardRoutes(rules: readonly ReadRule[] = READ_RULES): Hono<ApiEnv> {
const routes = new Hono<ApiEnv>();
for (const rule of rules) routes.on(rule.method, rule.path, readGuard(rule.capability));
return routes;
}
+974
View File
@@ -0,0 +1,974 @@
/**
* The quarterly calendar — a projection, not a table.
*
* Everything with a date on it already lives somewhere: contracts expire,
* obligations fall due, commitments open and close, holds lapse, export
* authorisations run out. This service reads those columns where they are and
* emits one common shape. Nothing here is stored, and nothing here can drift
* from the record it describes.
*
* Three things shape the implementation.
*
* **One query per source, each with its own date predicate and its own
* limit.** The convention elsewhere in this API is a flat `.limit(300)`
* ordered by `updated_at`, with the caller filtering by date in the browser —
* which means the deals actually closing this quarter are not guaranteed to be
* in the response at all. That is precisely the bug this endpoint exists to
* fix, so every predicate is server-side and every source is bounded
* independently rather than competing for one budget.
*
* **Totals are separate aggregate queries.** If the header counted the rows in
* the list it would under-report the moment any source truncated, and a
* quarterly figure that silently shrinks is worse than no figure. The counts
* are exact even when the list is cut short.
*
* **Renewal comes from `renewalAlarm()`.** The rule — expiry minus notice
* days, only when auto-renewal is on — is defined once, in the contracts
* service. The SQL below narrows candidates with the same arithmetic so the
* scan stays bounded, but every date and every state on an emitted event comes
* from calling that function. If the rule changes, it changes there.
*/
import {
and,
asc,
count,
eq,
gt,
gte,
isNotNull,
isNull,
lt,
or,
sql,
} from 'drizzle-orm';
import {
calendarEventId,
completableSpanState,
eventState,
quarterOf,
spanState,
type CalendarEvent,
type CalendarEventKind,
type Quarter,
} from '@pig/core';
import {
accounts,
allocations,
calendarEntries,
capacityCommitments,
complianceArtifacts,
contractObligations,
contracts,
demandDeals,
exportAuthorizations,
supplyDeals,
users,
type Database,
} from '@pig/db';
import { renewalAlarm } from './contracts';
/** Per-source ceiling. Generous enough that a real quarter never reaches it. */
const DEFAULT_SOURCE_LIMIT = 500;
export interface CalendarQuery {
from: Date;
/** Exclusive. Quarters are half-open so consecutive ones do not double-count. */
to: Date;
kinds?: readonly CalendarEventKind[];
accountId?: string;
ownerUserId?: string;
fiscalYearStartMonth?: number;
timeZone?: string;
sourceLimit?: number;
}
export interface CalendarTotals {
/**
* Σ acv × probability for deals whose expected close date falls in range.
* The number a GTM lead reads first, and nothing in PIG computed it before.
*/
weightedPipelineCents: number;
closingCount: number;
renewalCount: number;
obligationCount: number;
expiringAuthorizationCount: number;
}
export interface CalendarProjection {
from: string;
to: string;
quarter: Quarter;
events: CalendarEvent[];
/** True when any single source hit its limit; the totals are still exact. */
truncated: boolean;
totals: CalendarTotals;
}
/**
* Where the front end should go when an event is clicked.
*
* There is no record-detail route convention in this app yet — every page is
* flat — so the page is the load-bearing half and the query parameter is a
* hint the detail sheet can honour once one exists.
*/
function href(page: string, param: string, id: string): string {
return `/${page}?${param}=${id}`;
}
/** Drizzle returns numeric columns as strings; `probability` is one of them. */
function numeric(value: string | null): number | null {
if (value === null) return null;
const parsed = Number(value);
return Number.isFinite(parsed) ? parsed : null;
}
export class CalendarService {
constructor(
private readonly db: Database,
private readonly clock: () => Date = () => new Date(),
) {}
/**
* The reader's own quarter boundary.
*
* `users.timezone` is settable through PATCH /api/me/preferences and until
* now was read by nothing at all. A quarter is a local-midnight question, so
* this is the first place it genuinely matters — and UTC remains the honest
* fallback for a user who has never set one.
*/
async timeZoneFor(userId: string): Promise<string> {
const [row] = await this.db
.select({ timezone: users.timezone })
.from(users)
.where(eq(users.id, userId))
.limit(1);
return row?.timezone ?? 'UTC';
}
async project(query: CalendarQuery): Promise<CalendarProjection> {
const now = this.clock();
const timeZone = query.timeZone ?? 'UTC';
const fiscalYearStartMonth = query.fiscalYearStartMonth ?? 0;
const limit = query.sourceLimit ?? DEFAULT_SOURCE_LIMIT;
const wanted = query.kinds?.length ? new Set(query.kinds) : null;
const wants = (kind: CalendarEventKind): boolean => !wanted || wanted.has(kind);
const collected: { events: CalendarEvent[]; truncated: boolean }[] = await Promise.all([
wants('expected_close') ? this.expectedClose(query, now, limit) : empty(),
wants('contract_effective')
? this.contractDate(query, now, limit, 'contract_effective')
: empty(),
wants('contract_expiry')
? this.contractDate(query, now, limit, 'contract_expiry')
: empty(),
wants('contract_executed')
? this.contractDate(query, now, limit, 'contract_executed')
: empty(),
wants('renewal_notice') ? this.renewalNotices(query, now, limit) : empty(),
wants('obligation_due') ? this.obligations(query, now, limit) : empty(),
wants('capacity_window') ? this.capacityWindows(query, now, limit) : empty(),
wants('allocation_window') ? this.allocationWindows(query, now, limit) : empty(),
wants('hold_expiry') ? this.holdExpiries(query, now, limit) : empty(),
wants('supply_available_from') ? this.supplyAvailability(query, now, limit) : empty(),
wants('authorization_expiry') ? this.authorizationExpiries(query, now, limit) : empty(),
wants('artifact_expiry') ? this.artifactExpiries(query, now, limit) : empty(),
wants('calendar_entry') ? this.entries(query, now, limit) : empty(),
]);
const events = collected
.flatMap((source) => source.events)
.sort((a, b) => a.startsAt.localeCompare(b.startsAt) || a.id.localeCompare(b.id));
return {
from: query.from.toISOString(),
to: query.to.toISOString(),
quarter: quarterOf(query.from, fiscalYearStartMonth, timeZone),
events,
truncated: collected.some((source) => source.truncated),
totals: await this.totals(query),
};
}
// ------------------------------------------------------------------ totals
/**
* Counted in SQL rather than off the event list, so a truncated source
* cannot quietly shrink a quarterly figure. The kind filter is deliberately
* ignored here: narrowing the list to one kind should not blank the header
* the reader is narrowing against.
*/
private async totals(query: CalendarQuery): Promise<CalendarTotals> {
const { from, to, accountId, ownerUserId } = query;
const [pipeline, renewals, obligations, authorizations] = await Promise.all([
this.db
.select({
/**
* A closed-won deal forecasts at certainty and a closed-lost one at
* nothing, whatever `probability` still says; an open deal with no
* forecast contributes nothing rather than its full value, because
* an unfilled field is not a prediction of 100%.
*/
weightedCents: sql<string>`coalesce(sum(round(${demandDeals.acvCents} * (case
when ${demandDeals.stage} = 'closed_won' then 1
when ${demandDeals.stage} = 'closed_lost' then 0
else coalesce(${demandDeals.probability}, 0) end))), 0)`,
closing: sql<number>`count(*) filter (where ${demandDeals.stage} <> 'closed_lost')::int`,
})
.from(demandDeals)
.where(
and(
gte(demandDeals.expectedCloseDate, from),
lt(demandDeals.expectedCloseDate, to),
accountId ? eq(demandDeals.accountId, accountId) : undefined,
ownerUserId ? eq(demandDeals.ownerUserId, ownerUserId) : undefined,
),
),
this.db
.select({ value: count() })
.from(contracts)
.where(this.renewalPredicate(query)),
this.db
.select({ value: count() })
.from(contractObligations)
.innerJoin(contracts, eq(contracts.id, contractObligations.contractId))
.where(
and(
gte(contractObligations.dueAt, from),
lt(contractObligations.dueAt, to),
// Outstanding only. A count that includes work already done reads
// as a backlog that is not there.
isNull(contractObligations.completedAt),
accountId ? eq(contracts.accountId, accountId) : undefined,
ownerUserId ? eq(contractObligations.ownerUserId, ownerUserId) : undefined,
),
),
// An expiring authorisation has no owner column, so an owner filter can
// only ever exclude it — reporting zero rather than the whole book.
ownerUserId
? Promise.resolve([{ value: 0 }])
: this.db
.select({ value: count() })
.from(exportAuthorizations)
.where(
and(
gte(exportAuthorizations.expiresAt, from),
lt(exportAuthorizations.expiresAt, to),
accountId ? eq(exportAuthorizations.accountId, accountId) : undefined,
),
),
]);
return {
weightedPipelineCents: Math.round(Number(pipeline[0]?.weightedCents ?? 0)),
closingCount: pipeline[0]?.closing ?? 0,
renewalCount: renewals[0]?.value ?? 0,
obligationCount: obligations[0]?.value ?? 0,
expiringAuthorizationCount: authorizations[0]?.value ?? 0,
};
}
// ----------------------------------------------------------------- sources
private async expectedClose(query: CalendarQuery, now: Date, limit: number) {
const rows = await this.db
.select({ deal: demandDeals, accountName: accounts.name })
.from(demandDeals)
.leftJoin(accounts, eq(accounts.id, demandDeals.accountId))
.where(
and(
gte(demandDeals.expectedCloseDate, query.from),
lt(demandDeals.expectedCloseDate, query.to),
query.accountId ? eq(demandDeals.accountId, query.accountId) : undefined,
query.ownerUserId ? eq(demandDeals.ownerUserId, query.ownerUserId) : undefined,
),
)
.orderBy(asc(demandDeals.expectedCloseDate))
.limit(limit + 1);
return bounded(rows, limit, ({ deal, accountName }) => {
const at = deal.expectedCloseDate!;
const probability = numeric(deal.probability);
return {
id: calendarEventId('demand_deal', deal.id, 'expectedCloseDate'),
kind: 'expected_close' as const,
title: deal.name,
startsAt: at.toISOString(),
endsAt: null,
isSpan: false,
state: eventState({ at, now, completedAt: deal.closedAt }),
accountId: deal.accountId,
accountName,
ownerUserId: deal.ownerUserId,
amountCents: deal.acvCents,
currency: deal.currency,
recordType: 'demand_deal',
recordId: deal.id,
href: href('demand', 'deal', deal.id),
meta: {
stage: deal.stage,
probability,
productLine: deal.productLine,
weightedCents:
deal.acvCents !== null && probability !== null
? Math.round(deal.acvCents * probability)
: null,
},
};
});
}
private async contractDate(
query: CalendarQuery,
now: Date,
limit: number,
kind: 'contract_effective' | 'contract_expiry' | 'contract_executed',
) {
const column =
kind === 'contract_effective'
? contracts.effectiveAt
: kind === 'contract_expiry'
? contracts.expiresAt
: contracts.executedAt;
const field =
kind === 'contract_effective'
? 'effectiveAt'
: kind === 'contract_expiry'
? 'expiresAt'
: 'executedAt';
const label =
kind === 'contract_effective'
? 'takes effect'
: kind === 'contract_expiry'
? 'expires'
: 'executed';
const rows = await this.db
.select({ contract: contracts, accountName: accounts.name })
.from(contracts)
.leftJoin(accounts, eq(accounts.id, contracts.accountId))
.where(
and(
gte(column, query.from),
lt(column, query.to),
query.accountId ? eq(contracts.accountId, query.accountId) : undefined,
query.ownerUserId ? eq(contracts.ownerUserId, query.ownerUserId) : undefined,
),
)
.orderBy(asc(column))
.limit(limit + 1);
return bounded(rows, limit, ({ contract, accountName }) => {
const at = contract[field]!;
return {
id: calendarEventId('contract', contract.id, field),
kind,
title: `${contract.title} ${label}`,
startsAt: at.toISOString(),
endsAt: null,
isSpan: false,
// An executed date is a fact about the past, not an errand: it is
// recorded as done so it does not sit in the overdue list forever.
state:
kind === 'contract_executed'
? ('done' as const)
: eventState({ at, now, completedAt: contract.terminatedAt }),
accountId: contract.accountId,
accountName,
ownerUserId: contract.ownerUserId,
amountCents: contract.valueCents,
currency: contract.currency,
recordType: 'contract',
recordId: contract.id,
href: href('contracts', 'contract', contract.id),
meta: {
contractType: contract.type,
status: contract.status,
side: contract.side,
terminatedAt: contract.terminatedAt?.toISOString() ?? null,
},
};
});
}
/**
* The SQL narrows; `renewalAlarm()` decides.
*
* The predicate repeats the expiry-minus-notice arithmetic only to keep the
* scan bounded — the alternative is loading every auto-renewing contract in
* the book. Every date and state that reaches a caller comes from the shared
* function, so there is still exactly one definition of the rule.
*/
private renewalPredicate(query: CalendarQuery) {
return and(
eq(contracts.isAutoRenew, true),
isNotNull(contracts.noticeDays),
isNotNull(contracts.expiresAt),
// A terminated contract will not renew, so its notice date is not a
// deadline anyone should be chased about.
isNull(contracts.terminatedAt),
// The bounds are bound as ISO text and cast, not as `Date`: drizzle types
// parameters from the column in a comparison, and a raw template has no
// column to learn from, so postgres-js receives a Date it cannot encode
// and the whole request 500s. Found by calling the endpoint.
sql`${contracts.expiresAt} - make_interval(days => ${contracts.noticeDays}) >= ${query.from.toISOString()}::timestamptz`,
sql`${contracts.expiresAt} - make_interval(days => ${contracts.noticeDays}) < ${query.to.toISOString()}::timestamptz`,
query.accountId ? eq(contracts.accountId, query.accountId) : undefined,
query.ownerUserId ? eq(contracts.ownerUserId, query.ownerUserId) : undefined,
);
}
private async renewalNotices(query: CalendarQuery, now: Date, limit: number) {
const rows = await this.db
.select({ contract: contracts, accountName: accounts.name })
.from(contracts)
.leftJoin(accounts, eq(accounts.id, contracts.accountId))
.where(this.renewalPredicate(query))
.orderBy(asc(contracts.expiresAt))
.limit(limit + 1);
return bounded(rows, limit, ({ contract, accountName }) => {
const alarm = renewalAlarm(contract, now);
const at = alarm.renewalNoticeAt!;
return {
id: calendarEventId('contract', contract.id, 'renewalNoticeAt'),
kind: 'renewal_notice' as const,
title: `Renewal notice — ${contract.title}`,
startsAt: at.toISOString(),
endsAt: null,
isSpan: false,
// 'expired' means the window to give notice has gone; the notice date
// itself is simply late until then.
state:
alarm.renewalState === 'expired'
? ('overdue' as const)
: eventState({ at, now }),
accountId: contract.accountId,
accountName,
ownerUserId: contract.ownerUserId,
amountCents: contract.valueCents,
currency: contract.currency,
recordType: 'contract',
recordId: contract.id,
href: href('contracts', 'contract', contract.id),
meta: {
renewalState: alarm.renewalState,
expiresAt: contract.expiresAt?.toISOString() ?? null,
noticeDays: contract.noticeDays,
side: contract.side,
},
};
});
}
/**
* Every obligation on every contract, in one query.
*
* Obligations were reachable only inside GET /api/contracts/:id, so a
* quarter of them meant one request per contract. They are the dated things
* most likely to be missed, which makes that the wrong place for them to be.
*/
private async obligations(query: CalendarQuery, now: Date, limit: number) {
const rows = await this.db
.select({
obligation: contractObligations,
contract: contracts,
accountName: accounts.name,
})
.from(contractObligations)
.innerJoin(contracts, eq(contracts.id, contractObligations.contractId))
.leftJoin(accounts, eq(accounts.id, contracts.accountId))
.where(
and(
gte(contractObligations.dueAt, query.from),
lt(contractObligations.dueAt, query.to),
query.accountId ? eq(contracts.accountId, query.accountId) : undefined,
query.ownerUserId
? eq(contractObligations.ownerUserId, query.ownerUserId)
: undefined,
),
)
.orderBy(asc(contractObligations.dueAt))
.limit(limit + 1);
return bounded(rows, limit, ({ obligation, contract, accountName }) => ({
id: calendarEventId('contract_obligation', obligation.id, 'dueAt'),
kind: 'obligation_due' as const,
title: obligation.title,
startsAt: obligation.dueAt.toISOString(),
endsAt: null,
isSpan: false,
state: eventState({
at: obligation.dueAt,
now,
completedAt: obligation.completedAt,
}),
accountId: contract.accountId,
accountName,
ownerUserId: obligation.ownerUserId,
amountCents: null,
currency: null,
recordType: 'contract_obligation',
recordId: obligation.id,
href: href('contracts', 'contract', contract.id),
meta: {
obligationKind: obligation.kind,
contractId: contract.id,
contractTitle: contract.title,
completedAt: obligation.completedAt?.toISOString() ?? null,
},
}));
}
/**
* Commitment windows, split on the capacity shape where one is present.
*
* A commitment ramps and steps — it is not a rectangle — and `shape` is
* authoritative over `startsAt`/`endsAt` when set. Drawing one bar across
* the whole term shows a seller capacity in a month it does not exist in,
* which is exactly the mistake the shape column was added to prevent.
*/
private async capacityWindows(query: CalendarQuery, now: Date, limit: number) {
// No owner column anywhere on the supply chain of custody, so an owner
// filter cannot be satisfied and must exclude the source outright.
if (query.ownerUserId) return { events: [], truncated: false };
const rows = await this.db
.select({ commitment: capacityCommitments, accountName: accounts.name })
.from(capacityCommitments)
.leftJoin(accounts, eq(accounts.id, capacityCommitments.accountId))
.where(
and(
lt(capacityCommitments.startsAt, query.to),
gt(capacityCommitments.endsAt, query.from),
query.accountId ? eq(capacityCommitments.accountId, query.accountId) : undefined,
),
)
.orderBy(asc(capacityCommitments.startsAt))
.limit(limit + 1);
const truncated = rows.length > limit;
if (truncated) rows.length = limit;
const events: CalendarEvent[] = [];
for (const { commitment, accountName } of rows) {
const base = {
kind: 'capacity_window' as const,
isSpan: true,
accountId: commitment.accountId,
accountName,
ownerUserId: null,
amountCents: null,
currency: commitment.currency,
recordType: 'capacity_commitment',
recordId: commitment.id,
href: href('capacity', 'commitment', commitment.id),
};
const shape = commitment.shape;
const subSpans =
shape && shape.intervals.length >= 2 && shape.quantities.length >= 1
? shape.intervals.slice(0, -1).map((boundary, index) => ({
index,
startsAt: new Date(boundary),
endsAt: new Date(shape.intervals[index + 1]!),
gpuCount: shape.quantities[index] ?? commitment.gpuCount,
}))
: [
{
index: null,
startsAt: commitment.startsAt,
endsAt: commitment.endsAt,
gpuCount: commitment.gpuCount,
},
];
for (const span of subSpans) {
if (Number.isNaN(span.startsAt.getTime()) || Number.isNaN(span.endsAt.getTime())) {
continue;
}
if (span.startsAt >= query.to || span.endsAt <= query.from) continue;
events.push({
...base,
id: calendarEventId(
'capacity_commitment',
commitment.id,
span.index === null ? 'window' : `shape.${span.index}`,
),
title:
span.index === null
? commitment.name
: `${commitment.name}${span.gpuCount}× ${commitment.gpuType}`,
startsAt: span.startsAt.toISOString(),
endsAt: span.endsAt.toISOString(),
state: commitment.terminatedAt
? ('done' as const)
: spanState({ startsAt: span.startsAt, endsAt: span.endsAt, now }),
meta: {
gpuType: commitment.gpuType,
gpuCount: span.gpuCount,
envelopeGpuCount: commitment.gpuCount,
shaped: span.index !== null,
costPerGpuHourCents: commitment.costPerGpuHourCents,
terminatedAt: commitment.terminatedAt?.toISOString() ?? null,
},
});
}
}
return { events, truncated };
}
private async allocationWindows(query: CalendarQuery, now: Date, limit: number) {
if (query.ownerUserId) return { events: [], truncated: false };
const rows = await this.db
.select({
allocation: allocations,
commitmentName: capacityCommitments.name,
dealName: demandDeals.name,
accountId: demandDeals.accountId,
accountName: accounts.name,
})
.from(allocations)
.leftJoin(
capacityCommitments,
eq(capacityCommitments.id, allocations.capacityCommitmentId),
)
.leftJoin(demandDeals, eq(demandDeals.id, allocations.demandDealId))
.leftJoin(accounts, eq(accounts.id, demandDeals.accountId))
.where(
and(
lt(allocations.startsAt, query.to),
gt(allocations.endsAt, query.from),
query.accountId ? eq(demandDeals.accountId, query.accountId) : undefined,
),
)
.orderBy(asc(allocations.startsAt))
.limit(limit + 1);
return bounded(rows, limit, (row) => {
const { allocation } = row;
const gpuHours = numeric(allocation.gpuHours) ?? 0;
return {
id: calendarEventId('allocation', allocation.id, 'window'),
kind: 'allocation_window' as const,
title:
row.dealName ??
(allocation.internalTeam
? `Internal — ${allocation.internalTeam}`
: (row.commitmentName ?? 'Allocation')),
startsAt: allocation.startsAt.toISOString(),
endsAt: allocation.endsAt.toISOString(),
isSpan: true,
state:
allocation.releasedAt !== null
? ('done' as const)
: spanState({
startsAt: allocation.startsAt,
endsAt: allocation.endsAt,
now,
}),
accountId: row.accountId ?? null,
accountName: row.accountName ?? null,
ownerUserId: null,
// Revenue over the window, in cents — hours are fractional, money is not.
amountCents: Math.round(gpuHours * allocation.pricePerGpuHourCents),
currency: allocation.currency,
recordType: 'allocation',
recordId: allocation.id,
href: href('capacity', 'allocation', allocation.id),
meta: {
status: allocation.status,
guaranteeType: allocation.guaranteeType,
gpuHours,
internalTeam: allocation.internalTeam,
commitmentId: allocation.capacityCommitmentId,
releasedAt: allocation.releasedAt?.toISOString() ?? null,
},
};
});
}
/**
* A hold expiring is the one date on this calendar that changes what can be
* sold: the moment it passes, the capacity returns to everyone else's
* availability. It has never been visible anywhere.
*/
private async holdExpiries(query: CalendarQuery, now: Date, limit: number) {
if (query.ownerUserId) return { events: [], truncated: false };
const rows = await this.db
.select({
allocation: allocations,
dealName: demandDeals.name,
accountId: demandDeals.accountId,
accountName: accounts.name,
})
.from(allocations)
.leftJoin(demandDeals, eq(demandDeals.id, allocations.demandDealId))
.leftJoin(accounts, eq(accounts.id, demandDeals.accountId))
.where(
and(
gte(allocations.holdExpiresAt, query.from),
lt(allocations.holdExpiresAt, query.to),
query.accountId ? eq(demandDeals.accountId, query.accountId) : undefined,
),
)
.orderBy(asc(allocations.holdExpiresAt))
.limit(limit + 1);
return bounded(rows, limit, (row) => {
const at = row.allocation.holdExpiresAt!;
return {
id: calendarEventId('allocation', row.allocation.id, 'holdExpiresAt'),
kind: 'hold_expiry' as const,
title: `Hold expires — ${row.dealName ?? 'unassigned capacity'}`,
startsAt: at.toISOString(),
endsAt: null,
isSpan: false,
state: eventState({ at, now, completedAt: row.allocation.releasedAt }),
accountId: row.accountId ?? null,
accountName: row.accountName ?? null,
ownerUserId: null,
// What was turned away to keep the hold. Makes the deadline honest.
amountCents: row.allocation.holdOpportunityCostCents,
currency: row.allocation.currency,
recordType: 'allocation',
recordId: row.allocation.id,
href: href('capacity', 'allocation', row.allocation.id),
meta: {
status: row.allocation.status,
gpuHours: numeric(row.allocation.gpuHours),
commitmentId: row.allocation.capacityCommitmentId,
},
};
});
}
private async supplyAvailability(query: CalendarQuery, now: Date, limit: number) {
const rows = await this.db
.select({ deal: supplyDeals, accountName: accounts.name })
.from(supplyDeals)
.leftJoin(accounts, eq(accounts.id, supplyDeals.accountId))
.where(
and(
gte(supplyDeals.availableFrom, query.from),
lt(supplyDeals.availableFrom, query.to),
query.accountId ? eq(supplyDeals.accountId, query.accountId) : undefined,
query.ownerUserId ? eq(supplyDeals.ownerUserId, query.ownerUserId) : undefined,
),
)
.orderBy(asc(supplyDeals.availableFrom))
.limit(limit + 1);
return bounded(rows, limit, ({ deal, accountName }) => {
const at = deal.availableFrom!;
return {
id: calendarEventId('supply_deal', deal.id, 'availableFrom'),
kind: 'supply_available_from' as const,
title: `Capacity available — ${deal.name}`,
startsAt: at.toISOString(),
endsAt: null,
isSpan: false,
state: eventState({ at, now, completedAt: deal.closedAt }),
accountId: deal.accountId,
accountName,
ownerUserId: deal.ownerUserId,
amountCents: null,
currency: null,
recordType: 'supply_deal',
recordId: deal.id,
href: href('supply', 'deal', deal.id),
meta: {
stage: deal.stage,
gpuType: deal.gpuType,
gpuCount: deal.gpuCount,
targetCostPerGpuHourCents: deal.targetCostPerGpuHourCents,
},
};
});
}
/**
* An expired export authorisation silently converts lawful business into
* unlawful business. The schema says so and indexes the column for it, and
* until this endpoint nothing in PIG read it — no endpoint, no screen.
*/
private async authorizationExpiries(query: CalendarQuery, now: Date, limit: number) {
if (query.ownerUserId) return { events: [], truncated: false };
const rows = await this.db
.select({ authorization: exportAuthorizations, accountName: accounts.name })
.from(exportAuthorizations)
.leftJoin(accounts, eq(accounts.id, exportAuthorizations.accountId))
.where(
and(
gte(exportAuthorizations.expiresAt, query.from),
lt(exportAuthorizations.expiresAt, query.to),
query.accountId ? eq(exportAuthorizations.accountId, query.accountId) : undefined,
),
)
.orderBy(asc(exportAuthorizations.expiresAt))
.limit(limit + 1);
return bounded(rows, limit, ({ authorization, accountName }) => {
const at = authorization.expiresAt!;
return {
id: calendarEventId('export_authorization', authorization.id, 'expiresAt'),
kind: 'authorization_expiry' as const,
title: `Export authorisation expires — ${accountName ?? 'account'}`,
startsAt: at.toISOString(),
endsAt: null,
isSpan: false,
// Never 'done': an authorisation is not something anyone completes,
// and marking a lapsed one finished is the failure mode itself.
state: eventState({ at, now }),
accountId: authorization.accountId,
accountName,
ownerUserId: null,
amountCents: null,
currency: null,
recordType: 'export_authorization',
recordId: authorization.id,
href: href('accounts', 'account', authorization.accountId),
meta: {
authorizationType: authorization.authorizationType,
reference: authorization.reference,
// Rules in flux for this counterparty: re-verify, do not trust the date.
volatile: authorization.volatile,
evidenceUrl: authorization.evidenceUrl,
verifiedByUserId: authorization.verifiedByUserId,
},
};
});
}
private async artifactExpiries(query: CalendarQuery, now: Date, limit: number) {
if (query.ownerUserId) return { events: [], truncated: false };
const rows = await this.db
.select({ artifact: complianceArtifacts, accountName: accounts.name })
.from(complianceArtifacts)
.leftJoin(accounts, eq(accounts.id, complianceArtifacts.accountId))
.where(
and(
gte(complianceArtifacts.expiresAt, query.from),
lt(complianceArtifacts.expiresAt, query.to),
query.accountId ? eq(complianceArtifacts.accountId, query.accountId) : undefined,
),
)
.orderBy(asc(complianceArtifacts.expiresAt))
.limit(limit + 1);
return bounded(rows, limit, ({ artifact, accountName }) => {
const at = artifact.expiresAt!;
return {
id: calendarEventId('compliance_artifact', artifact.id, 'expiresAt'),
kind: 'artifact_expiry' as const,
title: `${artifact.claim} expires — ${accountName ?? 'account'}`,
startsAt: at.toISOString(),
endsAt: null,
isSpan: false,
state: eventState({ at, now }),
accountId: artifact.accountId,
accountName,
ownerUserId: null,
amountCents: null,
currency: null,
recordType: 'compliance_artifact',
recordId: artifact.id,
href: href('accounts', 'account', artifact.accountId),
meta: {
claim: artifact.claim,
scope: artifact.scope,
// Certification versus self-declared alignment decides procurement,
// so it travels with the deadline rather than being looked up later.
isCertified: artifact.isCertified,
soc2Type: artifact.soc2Type,
verifiedByUserId: artifact.verifiedByUserId,
},
};
});
}
private async entries(query: CalendarQuery, now: Date, limit: number) {
const rows = await this.db
.select({ entry: calendarEntries, accountName: accounts.name })
.from(calendarEntries)
.leftJoin(accounts, eq(accounts.id, calendarEntries.accountId))
.where(
and(
// A dated entry with no end is a point; one with an end is a span,
// and a span overlaps the window whenever it has not already closed.
lt(calendarEntries.startsAt, query.to),
or(
and(isNull(calendarEntries.endsAt), gte(calendarEntries.startsAt, query.from)),
and(isNotNull(calendarEntries.endsAt), gt(calendarEntries.endsAt, query.from)),
),
query.accountId ? eq(calendarEntries.accountId, query.accountId) : undefined,
query.ownerUserId ? eq(calendarEntries.ownerUserId, query.ownerUserId) : undefined,
),
)
.orderBy(asc(calendarEntries.startsAt))
.limit(limit + 1);
return bounded(rows, limit, ({ entry, accountName }) => ({
id: calendarEventId('calendar_entry', entry.id, 'startsAt'),
kind: 'calendar_entry' as const,
title: entry.title,
startsAt: entry.startsAt.toISOString(),
endsAt: entry.endsAt?.toISOString() ?? null,
isSpan: entry.endsAt !== null,
// Not `spanState`: this is the one projected row type with a completion
// column, so a closed window is overdue until `completed_at` says
// otherwise. Whether a missed QBR is flagged must not depend on whether
// its author happened to type an end time.
state: entry.endsAt
? completableSpanState({
startsAt: entry.startsAt,
endsAt: entry.endsAt,
now,
completedAt: entry.completedAt,
})
: eventState({ at: entry.startsAt, now, completedAt: entry.completedAt }),
accountId: entry.accountId,
accountName,
ownerUserId: entry.ownerUserId,
amountCents: null,
currency: null,
recordType: 'calendar_entry',
recordId: entry.id,
href: href('calendar', 'entry', entry.id),
meta: {
entryKind: entry.kind,
allDay: entry.allDay,
description: entry.description,
demandDealId: entry.demandDealId,
supplyDealId: entry.supplyDealId,
completedAt: entry.completedAt?.toISOString() ?? null,
},
}));
}
}
// ------------------------------------------------------------------- helpers
async function empty(): Promise<{ events: CalendarEvent[]; truncated: boolean }> {
return { events: [], truncated: false };
}
/**
* Each source asks for one row more than its budget. Detecting truncation any
* other way means either a second count query per source or silently returning
* a partial quarter as if it were whole.
*/
function bounded<Row>(
rows: Row[],
limit: number,
toEvent: (row: Row) => CalendarEvent,
): { events: CalendarEvent[]; truncated: boolean } {
const truncated = rows.length > limit;
if (truncated) rows.length = limit;
return { events: rows.map(toEvent), truncated };
}