Skip to content

Python API

Embed Landing with the same action contract and executor used by CLI and HTTP. The caller owns the database and runtime lifecycle; enter Runtime.running() while executing work.

Add Landing to your application's environment with uv add "landing==0.1.0". The isolated uv tool install path provides the CLI; embedding uses the package in your application's environment. Configure the model as described in Configuration.

Delegate an action

from pathlib import Path

from landing.models import ActionRequest
from landing.runtime import Runtime


async def review():
    async with Runtime(Path("landing.sqlite3")).running() as landing:
        return await landing.command("review", ActionRequest(
            mode="gatekeeper",
            instruction="Review the candidate against the acceptance criteria.",
            workspace=str(Path.cwd()),
            checks=["make acceptance"],
        ))

Commands select persisted modes: triage selects issuer, fix selects fixer, review selects gatekeeper, and explain selects explainer. Runtime.run() admits a request with its explicit mode. Read the returned action's status, decision, result, and error; exit_code() applies ordinary CLI semantics.

Pass workspaces={"candidate": Path("/srv/candidate")} to select registered names. create_app() accepts the same mapping, token, public origin, skills, and GitHub context for the HTTP service.

Streaming SDK

landing.agent.run_stream() accepts native Bub 0.5.0 options: session_id, text or content-part prompt, optional mutable state, per-call model, allowed_tools, allowed_skills, and reasoning_effort.

from contextlib import aclosing
from pathlib import Path

from landing.runtime import Runtime


async def explain():
    async with Runtime(Path("landing.sqlite3")).running() as landing:
        stream = await landing.agent.run_stream(
            session_id="release-question",
            prompt=',explain "Explain the failed release check."',
            allowed_tools=["fs.read", "bash", "skill"],
            allowed_skills=["release-investigation"],
        )
        async with aclosing(stream):
            async for event in stream:
                if event.kind == "text":
                    print(event.data["delta"], end="")
        return stream.error, stream.usage

Await the stream, consume it fully, and close it when leaving early. Closing unfinished work requests durable cancellation. Events are native text, reasoning, tool_call, tool_result, usage, error, and final; a final event ends a model step, not necessarily the whole task. Errors and usage remain on the stream. Validation and publication errors also persist in the action record. Native error kinds are retained; other execution failures use unknown.

Model failure logs retain available call metadata and validation error types without argument contents. These diagnostics do not establish the provider as the cause.

The four commands and mode are native agent tools. ,mode reads selection; ,mode gatekeeper selects it without creating a task. Selection persists in the workspace tape and is isolated by session. Content parts stay evidence rather than dispatching commands. Explicit state bypasses hook-based state loading. Per-call tools and skills only narrow mode limits; callers serialize turns within a session.

Hook integration

Pass an existing Bub framework to register Landing's business hooks alongside host hooks:

from pathlib import Path

from bub import BubFramework
from bub.channels.message import ChannelMessage
from landing.runtime import Runtime


async def handle():
    framework = BubFramework()
    framework.load_builtin_hooks()
    async with Runtime(Path("landing.sqlite3"), framework=framework).running() as landing:
        return await framework.process_inbound(ChannelMessage(
            session_id="release-question",
            channel="cli",
            content=',explain "Explain the failed release check."',
        ))

The message pipeline retains state, prompt, rendering, and dispatch hooks. Direct SDK calls return events without rendering or dispatching. Both paths share durable tasks and execution. The runtime binds the task workspace; host-provided native environments remain authoritative. Outbound channels belong to the host.

Skills and additional tools

Runtime(path, skill_dirs=[...]) and create_app(path, skill_dirs=[...]) add trusted roots with native discovery, skill loading, and $skill-name expansion. Configuration defines precedence.

Pass Bub Tool instances with Runtime(path, tools=[...]), then select them per mode. Authorization belongs in each tool and its execution environment.

Runtime design

All four modes share one Bub 0.5.0 agent loop. Landing registers hooks explicitly for mode state, prompts, a task sidecar, and execution. Standalone Landing does not discover external plugins or packaged channel skills.

The sidecar owns actions and action_events. SQLite tape storage in the same database reuses Bub's query and async adapter. Resetting model history does not remove task records; completed tasks do not replay. See Action records and Recovery.

Public objects

Bases: Model

Source code in src/landing/models.py
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
class ActionRequest(Model):
    mode: Mode
    instruction: str | None = None
    input: list[Input] = Field(default_factory=list)
    workspace: str | None = None
    checks: list[str] = Field(default_factory=list)

    @model_validator(mode="after")
    def validate_work(self) -> ActionRequest:
        if not (self.instruction and self.instruction.strip()) and not self.input:
            message = "Provide an instruction or input."
            raise ValueError(message)
        if any(not command.strip() for command in self.checks):
            message = "Check commands must not be empty."
            raise ValueError(message)
        if self.checks and self.mode not in {"fixer", "gatekeeper"}:
            message = "Check commands are supported by fixer and gatekeeper."
            raise ValueError(message)
        if len(self.model_dump_json().encode()) > MAX_REQUEST_BYTES:
            message = "The request exceeds 16 MiB."
            raise ValueError(message)
        return self

Bases: Model

Source code in src/landing/models.py
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
class Action(Model):
    id: str
    mode: Mode
    instruction: str | None
    workspace: str | None
    status: Status
    result: str | None = None
    decision: Decision | None = None
    error: dict[str, str] | None = None
    retry_of: str | None = None
    created_at: str
    updated_at: str
    started_at: str | None = None
    completed_at: str | None = None
    cancel_requested_at: str | None = None

    def exit_code(self) -> int:
        return int(self.status != "completed" or (self.mode == "gatekeeper" and self.decision != "allow"))
Source code in src/landing/runtime.py
 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
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
class Runtime:
    def __init__(
        self,
        path: Path,
        *,
        workspaces: Mapping[str, Path] | None = None,
        tools: Iterable[Tool] = (),
        skill_dirs: Iterable[Path] = (),
        framework: BubFramework | None = None,
        verify: Callable[[Action], None] | None = None,
    ) -> None:
        self.tasks = Tasks(path)
        self.store = SQLiteTapeStore(self.tasks.path)
        self.workspaces = workspaces
        self.verify = verify
        self.framework = framework or BubFramework(config_file=ConfigurationFile().config_file.expanduser())
        self.settings = ensure_config(Settings)
        if framework is None:
            self.framework.plugin_manager.register(SDKDefaults(self.framework), name="builtin")
        self.hooks = LandingHooks(self)
        self.framework.plugin_manager.register(self.hooks, name="landing")
        self.skill_dirs = tuple(Path(root).expanduser().resolve() for root in (*skill_dirs, *self.settings.skill_dirs))
        self.agent = Agent(
            self,
            tools=[*REGISTRY.values(), *TOOLS, DECIDE, NO_UPDATE, *tools],
            tape_store=self.store,
            skill_dirs=(),
        )
        self.active: dict[str, asyncio.Task[Action]] = {}
        self.execution = asyncio.Lock()

    def workspace(self, request: ActionRequest) -> Path:
        if self.workspaces is None:
            path = Path(request.workspace or ".").expanduser().resolve()
        else:
            name = request.workspace or "default"
            if name not in self.workspaces:
                message = f"Workspace {name!r} is not registered."
                raise ValueError(message)
            path = self.workspaces[name].resolve()
        if not path.is_dir():
            message = f"Workspace does not exist: {path}"
            raise ValueError(message)
        return path

    @contextlib.asynccontextmanager
    async def running(self) -> AsyncIterator["Runtime"]:
        # A process-wide worker owns this database. Read-only CLI queries need no lock.
        with self.tasks.path.open("rb") as owner:
            try:
                fcntl.flock(owner, fcntl.LOCK_EX | fcntl.LOCK_NB)
            except BlockingIOError as exc:
                message = "Another worker owns this database. Use --server or a different --db."
                self.store.close()
                self.tasks.close()
                raise ValueError(message) from exc
            try:
                self.tasks.recover()
                async with self.framework.running():
                    monitor = asyncio.create_task(self.watch_cancellations())
                    try:
                        yield self
                    finally:
                        monitor.cancel()
                        with contextlib.suppress(asyncio.CancelledError):
                            await monitor
                        await self.stop()
            finally:
                self.store.close()
                self.tasks.close()
                fcntl.flock(owner, fcntl.LOCK_UN)

    async def run(self, request: ActionRequest, *, retry_of: str | None = None) -> Action:
        self.workspace(request)
        action, _ = self.tasks.create(request, retry_of=retry_of)
        task = asyncio.create_task(self.execute(action.id))
        self.active[action.id] = task
        try:
            return await asyncio.shield(task)
        except asyncio.CancelledError:
            self.cancel(action.id)
            await asyncio.gather(task, return_exceptions=True)
            raise
        finally:
            self.active.pop(action.id, None)

    async def command(
        self, name: str, request: ActionRequest, *, session_id: str = "cli", scope: str = "cli", key: str | None = None
    ) -> Action:
        """Delegate a user action through its native command tool."""
        if name not in COMMANDS:
            message = f"Unknown command {name!r}. Choose triage, fix, review, or explain."
            raise ValueError(message)
        workspace = self.workspace(request)
        state = await self.framework.build_state({"_runtime_agent": self.agent.bub}, session_id)
        state.update(
            _runtime_workspace=str(workspace),
            landing_request=request.model_dump(),
            landing_scope=scope,
            landing_delivery_key=key,
        )
        prompt = shlex.join([self.agent.bub.command_prefix + name, "instruction=" + (request.instruction or "")])
        stream = await self.agent.run_stream(session_id=session_id, prompt=prompt, state=state)
        async with contextlib.aclosing(stream):
            async for _ in stream:
                pass
        return self.tasks.get(state["landing_action_id"])

    async def watch_cancellations(self) -> None:
        while True:
            for action_id, task in tuple(self.active.items()):
                if not task.done() and not task.cancelling() and self.tasks.get(action_id).cancel_requested_at:
                    task.cancel()
            await asyncio.sleep(0.1)

    async def checks(self, action_id: str, request: ActionRequest, workspace: Path) -> list[dict]:
        results = []
        for command in request.checks:
            shell = await shell_manager.start(cmd=command, cwd=str(workspace), session_id=action_id)
            timed_out = False
            try:
                async with asyncio.timeout(300):
                    await shell_manager.wait_closed(shell.shell_id)
            except TimeoutError:
                timed_out = True
                await shell_manager.terminate(shell.shell_id)
            item = {"command": command, "exit_code": shell.returncode, "output": shell.output, "timed_out": timed_out}
            self.tasks.event(action_id, "validation", item)
            results.append(item)
        return results

    async def consume(
        self,
        action_id: str,
        request: ActionRequest,
        workspace: Path,
        checks: list[dict],
        *,
        state: TurnState | None = None,
        events: asyncio.Queue | None = None,
        stream_state: StreamState | None = None,
    ) -> tuple[str, Decision | None]:
        row = self.tasks.connection.execute(
            "SELECT data FROM action_events WHERE action_id=? AND type='sdk.invocation' LIMIT 1", (action_id,)
        ).fetchone()
        invocation = json.loads(row[0]) if row else {}
        session_id = invocation.pop("session_id", action_id)
        supplied_prompt = invocation.pop("prompt", None)
        self.capabilities(request.mode, invocation)
        if state is None:
            state = await self.framework.build_state({"_runtime_agent": self.agent.bub}, session_id)
        state.update(landing_action_id=action_id, landing_mode=request.mode, _runtime_workspace=str(workspace))
        state.pop("landing_decision", None)
        state.pop("landing_llm_call", None)
        state.pop("landing_no_update", None)
        state.pop("landing_tool_failed", None)
        state.pop("allowed_skills", None)
        # Actions are serialized; discovery and the native skill tool share these per-turn SDK roots.
        self.agent.bub.skill_dirs = (workspace / ".agents/skills", *self.skill_dirs, Path.home() / ".agents/skills")
        # Content parts keep task evidence outside native command dispatch.
        stream = await self.agent.bub.run_stream(
            session_id=session_id,
            prompt=supplied_prompt if supplied_prompt is not None else task_prompt(request, checks),
            state=state,
            **invocation,
        )
        output = ""
        try:
            async with contextlib.aclosing(stream):
                async for event in stream:
                    if events is not None:
                        events.put_nowait(event)
                    if event.kind == "final" and "text" in event.data:
                        output = str(event.data["text"])
                        self.tasks.output(action_id, output)
        except Exception as exc:
            self.hooks.record_failure(action_id, state, exc)
            raise
        finally:
            if stream.error is not None:
                self.hooks.record_failure(action_id, state, stream.error)
            state.pop("landing_llm_call", None)
            if stream_state is not None:
                stream_state.error, stream_state.usage = stream.error, stream.usage
        if stream.error is not None:
            raise stream.error
        decision = self._decision(request, state, checks)
        if not output.strip():
            self.tasks.output(action_id, output, decision)
            message = "The model returned empty output."
            raise RuntimeError(message)
        self.hooks.record_completion(state)
        return output, decision

    def capabilities(self, mode, invocation) -> None:
        """Apply independent mode limits using the SDK's native tool resolution."""
        limits = self.settings.modes.get(mode, ModeSettings())
        if limits.allowed_tools is not None:
            available = self.agent.bub.tools
            configured = resolve_tool_names(limits.allowed_tools, all_names=available)
            requested = resolve_tool_names(invocation.get("allowed_tools"), all_names=available)
            invocation["allowed_tools"] = sorted(configured & requested)
        if limits.allowed_skills is not None:
            configured_skills = {name.casefold() for name in limits.allowed_skills}
            requested_skills = invocation.get("allowed_skills")
            invocation["allowed_skills"] = sorted(
                configured_skills
                if requested_skills is None
                else configured_skills & {name.casefold() for name in requested_skills}
            )

    @staticmethod
    def _decision(request: ActionRequest, state: TurnState, checks: list[dict]) -> Decision | None:
        if request.mode != "gatekeeper":
            return None
        if checks_failed(checks):
            return "block"
        return cast("Decision", state.get("landing_decision", "inconclusive"))

    async def perform(self, action_id: str, **kwargs) -> tuple[str, Decision | None]:
        request = self.tasks.request(action_id)
        workspace = self.workspace(request)
        checks = await self.checks(action_id, request, workspace) if request.mode == "gatekeeper" else []
        # Finish model-owned background processes before validating its changes.
        async with shell_manager.lifespan():
            output, decision = await self.consume(action_id, request, workspace, checks, **kwargs)
        if request.mode == "gatekeeper" and checks_failed(checks):
            decision = "block"
            output += "\nRequired validation failed; the change cannot proceed."
        self.tasks.output(action_id, output, decision)
        if request.mode == "fixer" and checks_failed(await self.checks(action_id, request, workspace)):
            message = "Required validation failed. Inspect the recorded checks and partial changes."
            raise RuntimeError(message)
        return output, decision

    async def execute(self, action_id: str, **kwargs) -> Action:
        async with self.execution:
            if not self.tasks.claim(action_id):
                return self.tasks.get(action_id)
            try:
                async with shell_manager.lifespan():
                    output, decision = await self.perform(action_id, **kwargs)
                if self.verify is not None:
                    self.verify(self.tasks.get(action_id))
            except asyncio.CancelledError:
                status = "cancelled" if self.tasks.get(action_id).cancel_requested_at else "interrupted"
                self.tasks.finish(action_id, status)
                raise
            except Exception as exc:
                return self.tasks.finish(
                    action_id,
                    "failed",
                    error={
                        "code": exc.kind.value if isinstance(exc, BubError) else ErrorKind.UNKNOWN.value,
                        "message": exc.message if isinstance(exc, BubError) else str(exc) or type(exc).__name__,
                    },
                )
            else:
                return self.tasks.finish(action_id, "completed", result=output, decision=decision)

    def stream(self, action_id: str, *, state: TurnState) -> AsyncStreamEvents:
        """Expose the shared executor's native events, with durable cancellation."""
        events: asyncio.Queue[StreamEvent | None] = asyncio.Queue()
        stream_state = StreamState()
        task: asyncio.Task | None = None

        async def drive():
            try:
                action = await self.execute(action_id, state=state, events=events, stream_state=stream_state)
                if action.error and stream_state.error is None:
                    stream_state.error = BubError(ErrorKind.UNKNOWN, action.error["message"])
                    events.put_nowait(StreamEvent("error", stream_state.error.as_dict()))
            finally:
                events.put_nowait(None)

        async def iterate():
            nonlocal task
            task = asyncio.create_task(drive())
            self.active[action_id] = task
            while (event := await events.get()) is not None:
                yield event
            await task

        async def close():
            if task is None or not task.done():
                self.cancel(action_id)
            if task is not None:
                await asyncio.gather(task, return_exceptions=True)
            self.active.pop(action_id, None)

        return AsyncStreamEvents(iterate(), state=stream_state, on_close=close)

    def cancel(self, action_id: str) -> Action:
        action = self.tasks.cancel(action_id)
        if (task := self.active.get(action_id)) and not task.done() and not task.cancelling():
            task.cancel()
        return action

    async def worker(self) -> None:
        try:
            while True:
                action_id = self.tasks.next()
                if action_id is None:
                    await asyncio.sleep(0.1)
                    continue
                task = asyncio.create_task(self.execute(action_id))
                self.active[action_id] = task
                try:
                    await task
                except asyncio.CancelledError:
                    if (parent := asyncio.current_task()) and parent.cancelling():
                        raise
                finally:
                    self.active.pop(action_id, None)
        finally:
            await self.stop()

    async def stop(self) -> None:
        for task in self.active.values():
            if not task.done():
                task.cancel()
        await asyncio.gather(*self.active.values(), return_exceptions=True)

running() async

Source code in src/landing/runtime.py
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
@contextlib.asynccontextmanager
async def running(self) -> AsyncIterator["Runtime"]:
    # A process-wide worker owns this database. Read-only CLI queries need no lock.
    with self.tasks.path.open("rb") as owner:
        try:
            fcntl.flock(owner, fcntl.LOCK_EX | fcntl.LOCK_NB)
        except BlockingIOError as exc:
            message = "Another worker owns this database. Use --server or a different --db."
            self.store.close()
            self.tasks.close()
            raise ValueError(message) from exc
        try:
            self.tasks.recover()
            async with self.framework.running():
                monitor = asyncio.create_task(self.watch_cancellations())
                try:
                    yield self
                finally:
                    monitor.cancel()
                    with contextlib.suppress(asyncio.CancelledError):
                        await monitor
                    await self.stop()
        finally:
            self.store.close()
            self.tasks.close()
            fcntl.flock(owner, fcntl.LOCK_UN)

run(request, *, retry_of=None) async

Source code in src/landing/runtime.py
137
138
139
140
141
142
143
144
145
146
147
148
149
async def run(self, request: ActionRequest, *, retry_of: str | None = None) -> Action:
    self.workspace(request)
    action, _ = self.tasks.create(request, retry_of=retry_of)
    task = asyncio.create_task(self.execute(action.id))
    self.active[action.id] = task
    try:
        return await asyncio.shield(task)
    except asyncio.CancelledError:
        self.cancel(action.id)
        await asyncio.gather(task, return_exceptions=True)
        raise
    finally:
        self.active.pop(action.id, None)

command(name, request, *, session_id='cli', scope='cli', key=None) async

Delegate a user action through its native command tool.

Source code in src/landing/runtime.py
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
async def command(
    self, name: str, request: ActionRequest, *, session_id: str = "cli", scope: str = "cli", key: str | None = None
) -> Action:
    """Delegate a user action through its native command tool."""
    if name not in COMMANDS:
        message = f"Unknown command {name!r}. Choose triage, fix, review, or explain."
        raise ValueError(message)
    workspace = self.workspace(request)
    state = await self.framework.build_state({"_runtime_agent": self.agent.bub}, session_id)
    state.update(
        _runtime_workspace=str(workspace),
        landing_request=request.model_dump(),
        landing_scope=scope,
        landing_delivery_key=key,
    )
    prompt = shlex.join([self.agent.bub.command_prefix + name, "instruction=" + (request.instruction or "")])
    stream = await self.agent.run_stream(session_id=session_id, prompt=prompt, state=state)
    async with contextlib.aclosing(stream):
        async for _ in stream:
            pass
    return self.tasks.get(state["landing_action_id"])
Source code in src/landing/server.py
 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
def create_app(  # noqa: C901 -- route definitions share an application lifespan.
    path: Path,
    *,
    workspaces: Mapping[str, Path] | None = None,
    token: str | None = None,
    base_url: str | None = None,
    github_repository: str | None = None,
    skill_dirs: Iterable[Path] = (),
) -> FastAPI:
    public_url = URL(base_url) if base_url else None
    if public_url and (
        public_url.scheme not in {"http", "https"}
        or not public_url.hostname
        or public_url.username is not None
        or public_url.query
        or public_url.fragment
        or public_url.path not in {"", "/"}
    ):
        message = "BASE_URL must be an HTTP(S) origin without credentials, a path, query, or fragment."
        raise ValueError(message)

    @asynccontextmanager
    async def lifespan(app: FastAPI):

        runtime = Runtime(
            path,
            workspaces=workspaces or {"default": Path.cwd()},
            skill_dirs=skill_dirs,
        )
        async with runtime.running():
            app.state.runtime = runtime
            worker = asyncio.create_task(runtime.worker())
            app.state.worker = worker
            try:
                yield
            finally:
                worker.cancel()
                with contextlib.suppress(asyncio.CancelledError):
                    await worker

    app = FastAPI(title="Landing", version=version("landing"), lifespan=lifespan, docs_url=None, redoc_url=None)

    @app.middleware("http")
    async def admission(request: Request, call_next):
        if request.url.path not in {"/healthz", "/up"}:
            if token and not hmac.compare_digest(request.headers.get("authorization", ""), "Bearer " + token):
                response = problem(401, "A valid bearer token is required.")
                response.headers["WWW-Authenticate"] = "Bearer"
                return response
            if (
                request.method == "POST"
                and request.url.path == "/v1/actions"
                and request.headers.get("content-type", "").split(";")[0] != "application/json"
            ):
                return problem(415, "Send an application/json request.")
            if len(await request.body()) > MAX_REQUEST_BYTES:
                return problem(413, "The request exceeds 16 MiB.")
        return await call_next(request)

    @app.exception_handler(RequestValidationError)
    async def validation_error(request, exc):
        status = 400 if any(error["type"] == "json_invalid" for error in exc.errors()) else 422
        return problem(status, "; ".join(error["msg"] for error in exc.errors()))

    @app.exception_handler(KeyError)
    async def not_found(request, exc):
        return problem(404, "The action was not found.")

    @app.exception_handler(ConflictError)
    async def conflict(request, exc):
        return problem(409, str(exc))

    @app.exception_handler(ValueError)
    async def invalid(request, exc):
        return problem(422, str(exc))

    @app.exception_handler(sqlite3.Error)
    async def storage_error(request, exc):
        return problem(503, "The database is unavailable.")

    def accept(runtime: Runtime, body: ActionRequest, key: str | None, retry_of: str | None = None) -> JSONResponse:
        if app.state.worker.done():
            return problem(503, "The worker is unavailable.")
        runtime.workspace(body)
        if github_repository and repository_context(github_repository) not in body.input:
            body = body.model_copy(update={"input": [*body.input, repository_context(github_repository)]})
        action, created = runtime.tasks.create(body, key=key, scope="server", retry_of=retry_of)
        return JSONResponse(
            action.model_dump(), status_code=201 if created else 200, headers={"Location": f"/v1/actions/{action.id}"}
        )

    def next_link(request: Request, **params) -> str:
        url = request.url.include_query_params(**params)
        if public_url:
            url = url.replace(scheme=public_url.scheme, netloc=public_url.netloc)
        return f'<{url}>; rel="next"'

    @app.post("/v1/actions", response_model=Action, status_code=201)
    async def create(
        body: ActionRequest,
        request: Request,
        idempotency_key: Annotated[str | None, Header(min_length=1, max_length=256)] = None,
    ):
        return accept(request.app.state.runtime, body, idempotency_key)

    @app.get("/v1/actions", response_model=list[Action])
    async def list_actions(
        request: Request, response: Response, limit: Annotated[int, Query(ge=1, le=100)] = 50, cursor: str | None = None
    ):
        tasks = request.app.state.runtime.tasks
        items = tasks.list(limit + 1, cursor)
        if len(items) > limit:
            response.headers["Link"] = next_link(request, cursor=items[limit - 1].id, limit=limit)
        return items[:limit]

    @app.get("/v1/actions/{action_id}", response_model=Action)
    async def view(action_id: str, request: Request):
        return request.app.state.runtime.tasks.get(action_id)

    @app.get("/v1/actions/{action_id}/events", response_model=list[Event])
    async def events(
        action_id: str,
        request: Request,
        response: Response,
        after: Annotated[int, Query(ge=0)] = 0,
        limit: Annotated[int, Query(ge=1, le=100)] = 50,
    ):
        items = request.app.state.runtime.tasks.events(action_id, after, limit + 1)
        if len(items) > limit:
            response.headers["Link"] = next_link(request, after=items[limit - 1].id, limit=limit)
        return items[:limit]

    @app.post("/v1/actions/{action_id}/cancellation", response_model=Action, status_code=202)
    async def cancel(action_id: str, request: Request):
        action = request.app.state.runtime.cancel(action_id)
        return JSONResponse(action.model_dump(), status_code=200 if action.status in TERMINAL else 202)

    @app.post("/v1/actions/{action_id}/retries", response_model=Action, status_code=201)
    async def retry(
        action_id: str,
        request: Request,
        idempotency_key: Annotated[str | None, Header(min_length=1, max_length=256)] = None,
    ):
        runtime = request.app.state.runtime
        return accept(runtime, runtime.tasks.request(action_id), idempotency_key, action_id)

    @app.get("/healthz")
    async def health():
        return {"status": "ok"}

    @app.get("/up", include_in_schema=False)
    async def ready():
        if app.state.worker.done():
            return problem(503, "The worker is unavailable.")
        app.state.runtime.tasks.connection.execute("SELECT id FROM actions LIMIT 1").fetchone()
        return {"status": "ok"}

    return app