Repository navigation
Expand file tree
/
Copy pathruntime.ts
More file actions
355 lines (334 loc) · 18.7 KB
/
Copy pathruntime.ts
File metadata and controls
355 lines (334 loc) · 18.7 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
/**
* The composition root for an agent run (#329 T9, epic #325).
*
* Every other module in `src/lib/agent/` takes what it needs as a parameter, which
* is what made them testable without a database, a model or a world. Something has
* to actually assemble those parts on a server, and this is the one place that
* does: the routes below `src/app/api/agent/` carry HTTP concerns only, and no
* route builds a tool context of its own.
*
* Three properties are deliberate:
*
* - **A run's connection is re-resolved from its persisted id, never from a
* request.** T2 forbids persisting a credential, so a run records a connection
* ID and nothing more; a drive therefore has to look the connection up again on
* the server. That is why an agent run requires a server-resolvable connection
* (`seed:…`): a connection defined only in a browser cannot be rebuilt by a
* process that is resuming somebody else's run, and taking one from the drive
* request's body would let a caller keep the id while pointing the run at a
* different server.
* - **The actor is read from the ledger, never from the caller.** `resolveConnection`
* is given the run's own persisted role, so a run resumed by anything at all still
* resolves exactly the connections its opener was allowed to see.
* - **Budget accounting and artifacts are process-wide and paired.** They are keyed
* by run id and released together when a run ends (`releaseExecutionRun`), so they
* have to outlive the request that started the run and be the SAME pair the next
* drive of that run sees.
*/
import { createDatabaseProvider } from "@/lib/db";
import { acquireExecutionProfileProvider } from "@/lib/db/factory";
import { editorExecutionContext } from "@/lib/api/execution-context";
import { type ExecutionArtifact, ExecutionArtifactStore } from "@/lib/db/operations/artifacts";
import { ExecutionBudgetTracker } from "@/lib/db/operations/budgets";
import { createCanonicalOperationRegistry } from "@/lib/db/operations/descriptors";
import { createTargetScope } from "@/lib/db/operations/policy";
import { LLMAuthError, LLMError, LLMRateLimitError } from "@/lib/llm/types";
import { type ExecutionProfileDenyCode, ExecutionProfileError } from "@/lib/db/errors";
import { logger } from "@/lib/logger";
import { resolveConnection, SeedConnectionError } from "@/lib/seed/resolve-connection";
import type { QueryResult } from "@/lib/types";
import { AgentRunDeadline } from "./deadline";
import { deriveDriveCeilings } from "./drive-budget";
import { AGENT_WORKFLOW_BUDGETS } from "./execution-policy";
import { type AgentInvestigationResult, runInvestigation } from "./investigation";
import { createAgentModel } from "./model-adapter";
import { AgentRepairLedger } from "./repair-ledger";
import { AgentRunService, AgentRunServiceError } from "./run-service";
import { AgentRunStore, resolveAgentLedgerWorld } from "./run-store";
import type { AgentRunFailureReason } from "./types";
/**
* How long a run's results stay readable. Sized against the run deadline rather than
* independently: an artifact that expired while its own run was still allowed to cite
* it would turn a verified claim into a dead reference, so the TTL is a comfortable
* multiple of the longest a run may live.
*
* The LONGEST deadline any workflow may take, not one workflow's — the store is
* process-wide and holds artifacts from every workflow at once, so a TTL derived from
* a shorter row would expire a `data-analysis` run's earliest evidence while
* that run was still going. A workflow with a shorter deadline simply gets more
* headroom than the multiple promises.
*/
const AGENT_ARTIFACT_TTL_MS =
Math.max(...Object.values(AGENT_WORKFLOW_BUDGETS).map((budget) => budget.runDeadlineMs)) * 4;
/**
* The entry cap of the process-wide artifact store, computed rather than picked:
* `45 statements × 4 runs = 180`. 45 is the largest per-workflow statement
* ceiling the frozen decision table admits (`database-assessment`); 4 is the
* number of runs this single agent process is assumed to carry at once.
*
* What that product bounds is FOUR DRIVES, not four runs, and the distinction was
* stated wrongly here until #373: this comment said "a run cannot produce more
* artifacts than it is allowed statements", which is not true of a run. A resumed
* drive's statement, elapsed-time and artifact ceilings are now derived from the
* run's own ledger rather than handed to it fresh (#999), so a run driven three
* times no longer triples its statement budget or evicts the evidence an earlier
* drive cited. One ceiling is still per drive — the repair ledger is rebuilt by
* each drive (`docs/BACKLOG.md` B6) — so repair attempts do start over on resume.
*
* The behaviour at the cap is worth knowing rather than guessing at. The cap is
* spent run-fairly (`ExecutionArtifactStore.put`): a store at the cap evicts the
* OLDEST ARTIFACT OF THE RUN THAT IS STORING, so a busy run cannot make "Show
* result" fail on a quieter one. The artifact allowance a resumed drive receives
* is derived the same way (#999), so a long-lived run no longer evicts its own
* earliest evidence while it is still live.
*
* Sized for the ceiling rather than for what a policy enforces at any one moment,
* so a statement budget lower than 45 leaves the cap correct and merely slack.
*/
export const AGENT_MAX_ARTIFACTS = 180;
/**
* The process-wide pair. Lazily built so that merely importing this module — which
* a route does at build time — allocates nothing while the runtime is off.
*/
let processResources: { tracker: ExecutionBudgetTracker; artifacts: ExecutionArtifactStore<QueryResult> } | null = null;
function runResources(): { tracker: ExecutionBudgetTracker; artifacts: ExecutionArtifactStore<QueryResult> } {
processResources ??= {
tracker: new ExecutionBudgetTracker(),
artifacts: new ExecutionArtifactStore<QueryResult>({
ttlMs: AGENT_ARTIFACT_TTL_MS,
maxArtifacts: AGENT_MAX_ARTIFACTS,
}),
};
return processResources;
}
/**
* One result this process still holds, or undefined (#329 T11).
*
* Undefined covers three states that are one answer to a caller: released with its
* run, expired by the TTL, or produced by a process that is not this one. All three
* mean the same thing — the rows are not here — and none of them is an error, which
* is why the route that surfaces this reports it as its own outcome rather than as a
* failure. `runId` is on the returned artifact so the caller can check the result it
* gets back belongs to the run it asked about; this store is process-wide.
*/
export function readAgentArtifact(correlationId: string, nowMs: number): ExecutionArtifact<QueryResult> | undefined {
return runResources().artifacts.get(correlationId, nowMs);
}
/**
* The run service, over the durable backend T1 selected.
*
* @throws AgentRunStoreError with `RUNTIME_DISABLED` while the agent flag is off,
* which is the default. Routes check the flag first and answer 404, so
* reaching this throw means the flag changed under a live request.
*/
export async function getAgentRunService(): Promise<AgentRunService> {
const world = await resolveAgentLedgerWorld();
return new AgentRunService({ store: new AgentRunStore({ world }), resources: runResources() });
}
/**
* Drives one run to its conclusion: start it if it is queued, resume it if a
* previous process left it running.
*
* Called two ways, and both must behave identically — in-process right after the
* run is opened, and from the authenticated drive callback when the durable
* transport asks for the run to be picked up again. Which one is driving is not
* something the run may observe: everything it needs is re-derived from the ledger.
*
* @throws AgentRunServiceError when the run does not exist or has already ended.
* @throws SeedConnectionError when the run's connection is no longer resolvable —
* it was removed, or the actor's role no longer reaches it.
*/
export async function driveAgentRun(runId: string): Promise<AgentInvestigationResult> {
const service = await getAgentRunService();
const report = await service.status(runId);
if (report === null) {
// Deliberately outside the recording below: a run that does not exist has no
// ledger to record a failure on, and creating one would manufacture the very
// record whose absence is being reported.
throw new AgentRunServiceError("RUN_NOT_FOUND", `agent run "${runId}" does not exist`);
}
try {
const { actor, connectionId } = report.record;
// The persisted actor is the sole authority: the role that decides which managed
// connections are visible is the one recorded when the run was opened.
const connection = await resolveConnection({ connectionId }, { role: actor.role, username: actor.sessionId });
// Capabilities and labels are type-driven and read without connecting, the way
// /api/db/provider-meta reads them. The live, read-only provider a statement
// actually runs on is acquired per call through the execution-profile seam.
// One provider for both, so a run's declared behaviour and its declared
// vocabulary can never come from two different readings.
//
// So no `withOneShotTunnel` wrapper here (#457): this provider is never connected.
// The one that runs statements comes from `acquireExecutionProfileProvider`, which
// opens the connection's pooled tunnel itself.
const provider = await createDatabaseProvider(connection);
const capabilities = provider.getCapabilities();
const labels = provider.getLabels();
// The ceilings a drive begins with are folded from the run's ledger, so a
// resumed drive inherits the spend its earlier drives recorded rather than
// starting each ceiling again.
const ceilings = deriveDriveCeilings(report.record, Date.now());
runResources().tracker.seedUsage(runId, {
executedStatements: ceilings.executedStatements,
totalElapsedMs: ceilings.executedMs,
});
runResources().artifacts.setRunAllowance(runId, ceilings.artifactAllowance);
return await runInvestigation(runId, {
service,
model: await createAgentModel(),
resources: {
connection,
capabilities,
labels,
registry: createCanonicalOperationRegistry(),
scope: createTargetScope(connectionId),
tracker: runResources().tracker,
artifacts: runResources().artifacts,
// The run's own workflow decides its wall clock, the same way it decides its
// statement budget and its turn ceiling. Read from the record the ledger
// returned, so a resumed drive is bounded by what the run was opened as.
deadline: new AgentRunDeadline(ceilings.deadlineMs),
repairs: new AgentRepairLedger(),
// Told the editor posture the run's persisted actor gets on this connection (non-admin DuckDB file access), which
// decides whether a SQLite handle opens at all, keys the profiled cache, and picks which open single-writer
// handle an operations acquisition may borrow.
acquireProvider: (target, profile) =>
acquireExecutionProfileProvider(target, profile, {}, editorExecutionContext(actor, target)),
},
});
} catch (error) {
// A drive refused because another already owns the run is not a failure OF the
// run: it is healthy and in flight. Recording `failed` here would end the very
// run the first drive is still carrying, so the refusal is left to the caller
// (the drive route answers 409) and nothing is written to the ledger.
if (!(error instanceof AgentRunServiceError && error.reasonCode === "RUN_ALREADY_DRIVEN")) {
await recordDriveFailure(service, runId, error);
}
throw error;
}
}
/**
* Writes the ending a dead drive owes the run, without ever replacing the reason
* the caller needs to see.
*
* `runInvestigation` ends a run it entered, so this covers the window before and
* around it: resolving the connection, reading capabilities, building the model.
* A throw there used to unwind past the ledger completely, leaving a run at
* `queued` with an empty timeline whose reason existed only in the server log —
* and with no drive producer yet (`docs/BACKLOG.md` B9), nothing would return to it.
*
* Every failure of the recording itself is swallowed, on purpose. The run may have
* ended between the throw and this call, or have an execution still in flight; both
* make `finish` throw, and neither is what the caller asked about. Losing the
* original error to a bookkeeping error would trade a diagnosable failure for a
* confusing one.
*/
async function recordDriveFailure(service: AgentRunService, runId: string, error: unknown): Promise<void> {
const reason = classifyDriveFailure(error);
try {
await service.finish(runId, "failed", { reason });
} catch (recordingError) {
logger.error("Agent run failed and its ending could not be recorded", recordingError, { runId, reason });
}
}
/**
* The profile refusals that are about the connection's agent credential rather than
* about the engine. Kept beside the classifier that reads them: the set exists only to
* split one error type into the two things a user can do about it.
*/
const AGENT_CREDENTIAL_DENY_CODES: ReadonlySet<ExecutionProfileDenyCode> = new Set([
"AGENT_CREDENTIAL_UNRESOLVABLE",
"AGENT_CREDENTIAL_WITH_CONNECTION_STRING",
]);
/**
* The profile refusal that is about the database PRINCIPAL rather than about the
* engine or the credential field that names it.
*
* Kept beside the other set for the same reason: one error type, three things a user
* can do about it. This one is raised when the credential resolved and the provider
* exists, and the profile then refused the user it opened as, so neither of the other
* two labels is true of it.
*/
const PROFILE_PRINCIPAL_DENY_CODES: ReadonlySet<ExecutionProfileDenyCode> = new Set<ExecutionProfileDenyCode>([
"PROFILE_PRIVILEGES_TOO_BROAD",
"PROFILE_PRIVILEGES_TOO_NARROW",
]);
/**
* Chooses the label a user sees from the error's TYPE.
*
* Never from its message: that text comes from a model provider, a driver or a
* connection resolver, and none of them promise to keep a key, a host name or an
* internal path out of it. The message goes to the log; only the label crosses to
* the browser.
*/
function classifyDriveFailure(error: unknown): AgentRunFailureReason {
// `instanceof` rather than the marker checks `mapAgentModelError` needs: both of
// these are this repository's own classes, so there is one copy of each in the
// bundle. The SDK's errors are the ones that arrive in duplicate.
if (error instanceof SeedConnectionError) return "connection-unresolvable";
// An engine with no read-only profile is not a server fault, and calling it
// `internal` sent a user to a log that could only tell them what their own
// connection already says. Checked before the generic classes below because it is
// the specific thing that happened.
//
// The reason code, not just the type: `resolveAgentCredential` raises this same
// error on ANY engine for a credential that cannot be applied, and calling that
// "engine unsupported" told a PostgreSQL operator something false about their
// database while saying nothing about the credential they could fix (B47).
//
// The principal codes are the third cause, split off the same way (the SQL Server
// work of 2026-09-18). They say the engine granted the profile and then refused the
// USER: `sa`, which is the only SQL Server credential this repository's own
// `database-compose.yml` ships, is refused by the profile while `libredb_agent`
// beside it is accepted. Reported as "engine unsupported", that told an operator to
// change engines when the fix is one CREATE LOGIN.
//
// TWO codes and one label, deliberately. A principal can be refused for holding too
// much or for holding too little - SQL Server's admission step needs `SHOWPLAN`, so a
// plain reader cannot be admitted - and the REPAIRS are opposite, which is why they are
// separate codes carrying opposite advice (`PROFILE_REFUSAL_ADVICE`). What a RUN ended
// as is the same fact either way: this connection's user is not one the profile will
// run as, which is what this label says and all it says.
if (error instanceof ExecutionProfileError) {
if (AGENT_CREDENTIAL_DENY_CODES.has(error.reasonCode)) return "agent-credential-unusable";
if (PROFILE_PRINCIPAL_DENY_CODES.has(error.reasonCode)) return "agent-principal-refused";
return "engine-unsupported";
}
/*
These were one label until 2026-08-12, on the reasoning that the user's next move
is the same for all of them — look at the model settings. A live run falsified it.
A Gemini free-tier quota (15 requests a minute) was exhausted by testing, and the
rail told the user the provider "is not configured or could not be reached" about a
provider that was configured and had answered seconds earlier. For a quota the next
move is to wait a minute; for a refused key it is to fix the key; only the rest send
anyone to the settings. The subclasses are checked before `LLMError` itself, which
all of them extend.
*/
if (error instanceof LLMRateLimitError) return "model-rate-limited";
if (error instanceof LLMAuthError) return "model-unauthorized";
if (error instanceof LLMError) {
/*
A run collapsed to the widest reason still says why, in the log if not in the ledger.
`model-unavailable` is the catch-all of the three and carries nothing about WHAT the provider
did: the ledger records the verdict and never the error behind it. Measured 2026-09-21,
`ministral-3:3b` lost eight runs here, every one of them AFTER five to ten successful tool
calls, and `isRetryableError` retried none - so the fault was not in the retryable class and
widening that bound would have changed nothing. Which subclass and message it is was the only
thing that could say what to fix, and there was no way to find out.
An operator reading a run that died "model-unavailable" is in the same position. So the error
is named at warn level - its class and the first 200 characters, which is what distinguishes a
refused connection from a context overflow from a model that was pulled mid-run - while the
reason returned to the run is unchanged.
*/
logger.warn("LLMError mapped to model-unavailable", {
route: "agent/runtime",
error: error.name,
message: error.message.slice(0, 200),
});
return "model-unavailable";
}
// Everything else, deliberately unnamed: a provider that could not be built and a
// programming error are both "this server could not carry the run", and guessing
// between them would put a claim on the rail that this function cannot support.
return "internal";
}