Compare commits

..
Author SHA1 Message Date
sysadmin 108cbfa173 Merge pull request 'feat(controller): expose instance-level fleet identity and health snapshots (Closes #978)' (#979) from feat/issue-978-instance-fleet-snapshot into master 2026-07-30 09:34:33 -05:00
sysadminandClaude Opus 4.8 0dbff6dcd5 fix(controller): production client_instance_id launch path and instance-aware mutation gate
Remediate review 655 blockers on PR #979 (issue #978):

B1 — Production launcher (mcp_application_launcher) mints one trusted
client_instance_id per application launch and propagates it to all five
namespace workers via GITEA_MCP_CLIENT_INSTANCE. launcher_entry and
multi_namespace_launcher_entries use that path. Workers never invent a
trusted ID; missing/malformed/user-supplied values fail closed.

B2 — _check_mcp_runtimes_diagnostics is instance-aware: two legitimate
instances sharing a profile are allowed when each has a distinct trusted
client_instance_id; duplicate workers for the same (instance, profile)
still fail closed. Worker identity/generation exported for peer scans.

Tests cover shared ID across five namespaces, distinct launches, multi-
instance same profile, same-instance duplicates, untrusted attribution,
and the production launcher path.

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
2026-07-30 10:12:16 -04:00
sysadminandClaude Opus 4.8 4a1cc63e94 feat(controller): expose instance-level fleet identity and health snapshots
Add a pure fleet snapshot assessor and a read-only controller/reconciler
MCP tool so multi-instance fleets can be enumerated by client_instance_id
without treating shared client_type as duplication. Trusted launcher
instance IDs, classification (missing/unmanifested/collisions/historical),
registry schema extensions, docs, and regression tests for #978.

Closes #978

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
2026-07-30 03:14:15 -04:00
sysadmin 6596b259fc Merge pull request 'fix(runtime): recognize client identity environment and refresh worker registrations' (#976) from fix/issue-975-client-identity-heartbeat into master 2026-07-29 23:26:34 -05:00
sysadminandClaude Opus 4.8 b0868be6b3 fix(runtime): recognize heartbeat interval env and wire production lifecycle tests
Address REQUEST_CHANGES review #652 on PR #976 (issue #975).

Blocker 1: name GITEA_WORKER_HEARTBEAT_INTERVAL_SECONDS individually in
RECOGNIZED_GITEA_ENV_KEYS so the production-consumed operator override is no
longer classified as unsupported-env / runtime_reconnect_required. No prefix
broadening; unknown overrides remain fail-closed.

Blocker 2: add focused production-path tests that call _active_worker_identity
and _start_worker_heartbeat, proving supervisor attachment, identity/session/
generation/pid/epoch agreement, no duplicates, failed/fenced paths, orderly
shutdown, configured interval, and off-loop sqlite heartbeats.

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
2026-07-29 23:41:37 -04:00
jcwalker3andClaude Opus 4.8 1a38ef95e3 fix(runtime): recognize client identity environment and refresh worker registrations
Two defects left behind by #948 made sanctioned multi-client operation
impossible. They are inseparable: fixing either alone still leaves the
multi-client canary unable to run.

1. Client-identity environment keys were not recognized.

   gitea_mcp_server reads GITEA_MCP_CLIENT, GITEA_MCP_CLIENT_INSTANCE and
   GITEA_MCP_CLIENT_SESSION as the authoritative client-identity inputs for
   worker registration, but none of the three appeared in
   RECOGNIZED_GITEA_ENV_KEYS or matched a recognized prefix. The runtime
   diagnostic scans every peer mcp_server.py process environment and classifies
   any unlisted GITEA_* key as an unsupported override, which is raised as a
   runtime blocker, so gitea_resolve_task_capability returned
   blocker_kind=runtime_reconnect_required with stop_required=true. Because that
   resolver is the mandatory preflight for every author, reviewer and merger
   mutation, setting the very variable #948 requires closed the mutation gate
   for the whole fleet, and reconnecting could not clear it: the variable is
   re-exported from the client's server definition on every launch.

   The three keys are now named individually in the recognized-key set. No
   prefix is added, so an unrecognized GITEA_* override is still refused
   exactly as before.

   A related inconsistency in the same path is also fixed. The diagnostic
   reasons are raised as one RuntimeError, but the preflight re-raise
   recognized only "stale-runtime:", so an "unsupported-env:" reason was
   silently swallowed there while still failing the resolver. Both reason
   families now live in RUNTIME_DIAGNOSTIC_HARD_PREFIXES beside the function
   that produces them, and both propagate identically. This only widens what is
   refused, never what is permitted.

2. WorkerRegistry.heartbeat() had no production caller.

   #948 delivered heartbeat() but only tests called it. The single production
   writer registers once per process behind an attempted-once flag, and
   register() stamps the same timestamp into both started_at and
   last_heartbeat_at. Nothing advanced it afterwards: no lifespan hook, no
   background task, no atexit handler in a process that blocks in mcp.run().
   Since liveness is age against heartbeat_ttl_seconds, that TTL was not a
   liveness window at all but a hard cap on how long any client could stay
   attached; at 900 seconds a healthy, connected, client-managed process became
   session_ownership=unowned with blocker_kind=session_attachment_missing.

   WorkerHeartbeatSupervisor in mcp_worker_identity is the missing caller,
   started from _active_worker_identity() at the moment register() succeeds,
   because that is the only point where identity and fencing_epoch are both
   known. It is a daemon thread rather than an asyncio task or a request-driven
   refresh because renewal must survive an idle session, and because the
   registry performs blocking BEGIN IMMEDIATE sqlite writes that must not run on
   the server's event loop. daemon=True is deliberate: a hard kill takes the
   thread with it, so a dead worker still goes stale on the normal TTL.

   heartbeat_interval_for() returns one third of the TTL, hard-capped at one
   half, so two consecutive beats can be lost without the row expiring and no
   override can produce an interval that outlives the registration it renews.
   heartbeat() gains optional keyword-only expectations (session, generation,
   client name, pid); each supplied one must match the recorded row or the
   renewal is refused with the existing BLOCKER_FENCED literal rather than a new
   blocker_kind, since consumers switch on that value. Omitting them preserves
   the pre-existing behavior exactly. Client names are compared normalized, so
   several namespaces of one application stay one client while separate
   applications stay distinct.

   A terminal refusal stops the supervisor permanently and records why, so a
   fenced session can never beat its way back into ownership. A transient
   failure is counted and beating continues. An atexit hook stops it on orderly
   shutdown. status() is surfaced read-only as worker_heartbeat on
   gitea_get_runtime_context so a stopped heartbeat is diagnosable before the
   TTL turns it into session_attachment_missing; it grants nothing.

   claim_generation() still has no production caller. It bumps fencing_epoch,
   which would fence the supervisor's cached epoch, and the strict refusal is
   left in place deliberately: auto-re-adopting a bumped epoch would defeat
   fencing.

No lock or lease TTL is changed, including the author issue-lock TTL, and no
mutation refused today becomes permitted.

Tests: tests/test_issue_975_client_identity_heartbeat.py adds 40 focused tests
covering all 13 acceptance criteria. Every TTL assertion uses an injected
clock; no test waits for a real TTL. The thread-loop tests use a
millisecond-scale interval with bounded polling.

Focused: 40 passed, 15 subtests passed.
Full suite from inside the branches worktree: 28 failed, 6105 passed, 6 skipped,
1105 subtests — an identical failure set to the 324a0c8a baseline measured in a
sibling branches worktree (28 failed, 6065 passed, 1090 subtests). Zero new
failures; the delta is exactly the added tests.

Closes #975

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
2026-07-29 20:06:10 -05:00
sysadmin 324a0c8a8d Merge pull request 'fix(runtime): support validated cross-repository canonical roots (#973)' (#974) from fix/issue-973-cross-repo-canonical-roots into master 2026-07-29 17:45:44 -05:00
11 changed files with 5106 additions and 61 deletions
+198
View File
@@ -0,0 +1,198 @@
# Instance-level fleet identity and health snapshots
Issue **#978**. Companion primitives: **#948** (worker ownership), **#975**
(heartbeat lifecycle), **#951** (restart receipts).
## Why this exists
The multi-instance fleet gate (#963) needs a *native* answer to:
* which application launches are live;
* which of the five namespace workers belong to each launch;
* whether two Codex (or Claude, Grok, …) launches are distinct;
* whether any live collision is real rather than a shared client type.
Process tables, PID proximity, configuration files, and direct SQLite access
are **not** production evidence. The sanctioned surface is the read-only MCP
tool `gitea_snapshot_instance_fleet` on **controller** and **reconciler**
namespaces.
## Identity hierarchy
| Identity | Scope | Who generates it | Lifetime |
| --- | --- | --- | --- |
| `client_type` | Application family (`codex`, `claude_code`, `gemini`, `grok`, …) | Launcher sets `GITEA_MCP_CLIENT` | Stable for the product |
| `client_instance_id` | One running application launch | **Trusted launcher**, once per launch, as `GITEA_MCP_CLIENT_INSTANCE` | Fresh launch → new ID; reconnect of same launch → same ID; full app restart → new ID |
| `fleet_run_id` | Operator-approved enrollment / canary cohort | Operator / controller sets `GITEA_MCP_FLEET_RUN_ID` | Duration of the approved rollout |
| `namespace` | `author` \| `reviewer` \| `merger` \| `controller` \| `reconciler` | Profile / MCP server binding | Process lifetime |
| `worker_id` / `worker_identity` | One namespace worker process | Runtime registry at first registration | New on worker restart; not reused |
| `session_id` | Client session ownership | Launcher `GITEA_MCP_CLIENT_SESSION` or runtime | Session lifetime |
| `generation_id` | One daemon launch | Runtime at process boot | New on process restart |
| `process_identity` / PID | OS process | Runtime | Process lifetime |
### Rules
1. Multiple active instances **may** share the same `client_type`.
2. Every application launch receives a **distinct** `client_instance_id`.
3. All five namespace workers of one launch report the **same**
`client_instance_id`.
4. Each namespace worker has a **distinct** `worker_identity`, process
identity, generation, and PID.
5. Instance identity is **never** inferred from PID proximity, timestamps, or
client type alone.
6. Live reuse of one `client_instance_id` with conflicting workers fails closed.
7. Sharing only a profile or `client_type` is **not** a duplicate.
This deliberately **replaces** any permanent `exactly_one_per_profile` fleet
model (#949 assumption) as the operating rule for multi-instance fleets.
## How five workers join one instance
1. The host starts one application instance (for example one Codex session).
2. The **production application launcher**
(`mcp_application_launcher.build_application_mcp_servers` /
`gitea_config.multi_namespace_launcher_entries`) mints exactly one
`client_instance_id` via `mcp_fleet_snapshot.generate_client_instance_id`
and injects the same env into every namespace worker:
```text
GITEA_MCP_CLIENT=<client_type>
GITEA_MCP_CLIENT_INSTANCE=<client_instance_id>
GITEA_MCP_INSTANCE_PROVENANCE=trusted_launcher
GITEA_MCP_CLIENT_SESSION=<session_id> # optional but recommended
GITEA_MCP_FLEET_RUN_ID=<enrollment id> # when on an approved canary
GITEA_CLIENT_MANAGED=1
GITEA_MCP_PROFILE=<role-profile>
```
3. The MCP client starts the five namespace processes (`gitea-author`,
`gitea-reviewer`, `gitea-merger`, `gitea-controller`, `gitea-reconciler`)
from that generated `mcpServers` map — each entry carries the **same**
instance ID.
4. Each worker registers once into the worker registry with its own
`worker_identity`, `namespace`, `generation_id`, and PID, but the shared
`client_instance_id`. Workers never mint a trusted instance ID themselves.
Only IDs matching the launcher format `inst-<client>-<timestamp>-<digest>`
are trusted. Missing, `legacy-pid-*` / `pid-*` placeholders, and ordinary
user-supplied strings fail closed as untrusted and cannot authorize
multi-instance fleet mutation safety.
Worker reconnect (same process, same registration) keeps the instance ID.
Worker restart (new process) mints a new worker identity and generation but
must still receive the same `GITEA_MCP_CLIENT_INSTANCE` from the parent
application if it is the same launch. Full application restart mints a new
`client_instance_id` (call the launcher without a prior ID).
### Mutation gate (instance-aware)
The capability / runtime diagnostic gate
(`_check_mcp_runtimes_diagnostics`) is **instance-aware**:
* Two legitimate application instances that share a profile (same
`GITEA_MCP_PROFILE`) are allowed when each has a distinct trusted
`client_instance_id`.
* More than one live worker for the same
`(client_instance_id, profile/namespace)` fails closed as a duplicate
namespace worker.
* Multiple processes sharing a profile **without** trusted instance
evidence still fail closed (indistinguishable from a duplicate).
## Approved fleet manifest
An operator (or controller enrollment step) obtains an approved manifest as a
list of expected instances, for example:
```json
{
"instances": [
{
"client_type": "codex",
"client_instance_id": "inst-codex-20260730T120000Z-abc123def456",
"fleet_run_id": "canary-963-2026-07-30",
"namespaces": ["author", "reviewer", "merger", "controller", "reconciler"]
},
{
"client_type": "codex",
"client_instance_id": "inst-codex-20260730T120100Z-fed654cba321",
"fleet_run_id": "canary-963-2026-07-30",
"namespaces": ["author", "reviewer", "merger", "controller", "reconciler"]
}
]
}
```
Pass that list as `expected_manifest` to `gitea_snapshot_instance_fleet`.
When a manifest is supplied:
* instances on the manifest but not live → `missing_expected`;
* live instances not on the manifest → `unmanifested`;
* both are **active blockers** for fleet-gate safety.
Without a manifest, the snapshot still enumerates the live fleet and classifies
identity collisions; it does not invent enrollment policy.
## Snapshot consistency
Each snapshot includes:
* `snapshot_at` — UTC timestamp;
* `consistency_token` / `registry_revision` — content digest over worker
identity, instance id, status, heartbeat, fencing, and generation.
Two successive snapshots can prove heartbeat continuity via
`mcp_fleet_snapshot.compare_snapshot_heartbeats`.
## Classification (active vs historical)
**Active blockers** (make `live_fleet_safe=false`):
* missing expected instance;
* unmanifested extra instance;
* duplicate namespace worker within one instance;
* live `client_instance_id` collision;
* reused worker / session / generation / process / PID / fencing identity;
* orphaned or unowned workers;
* unknown client;
* foreign-repository workers;
* old-revision workers;
* stale workers still marked active;
* legacy incomplete instance identity.
**Historical dead rows** (`status` released/superseded) are reported as
`historical_dead` findings with `active_blocker=false`. They never
automatically make the live fleet unsafe.
## Fail-closed behaviour and recovery
| Condition | Mutation safety | Diagnostic reads | Recovery |
| --- | --- | --- | --- |
| Two instances, same `client_type`, distinct IDs | Safe (if otherwise healthy) | Available | None needed |
| Live reuse of one `client_instance_id` | Unsafe | Available | Stop the colliding launch or re-issue a distinct ID |
| Two workers, same namespace, one instance | Unsafe | Available | Stop the extra worker |
| Legacy registration without trusted instance ID | Unsafe for fleet mutation | Available | Relaunch with launcher-issued `GITEA_MCP_CLIENT_INSTANCE` |
| Historical dead row only | Does not block alone | Available | No action required for fleet safety |
| Registry unavailable | Snapshot fails closed | N/A | Repair registry path / reconnect namespaces |
Incomplete or untrusted instance identity **cannot** authorize unsafe
mutation. Diagnostic reads remain available where `gitea.read` allows.
## Permissions
* **Allowed:** `controller`, `reconciler` with `gitea.read`.
* **Denied:** author, reviewer, merger (even with `gitea.read` for other tools).
* **No new mutation permissions** are granted to any role.
## Related surfaces
* `mcp_worker_identity` — registry, worker identity, heartbeats (#948, #975).
* `mcp_fleet_snapshot` — pure snapshot + classification (#978).
* `gitea_snapshot_instance_fleet` — sanctioned MCP tool (#978).
* `gitea_get_runtime_context` — single-process view (not fleet-wide).
## Non-goals
* Starting the #963 canary.
* Purging historical registry rows.
* Preserving “exactly one process per profile” as the permanent model.
* Using shell / process-table / SQLite inspection as production fleet evidence.
+1
View File
@@ -159,6 +159,7 @@ that gates each call, not which tools exist.
- `gitea_sentry_reconcile_issue`
- `gitea_sentry_watchdog`
- `gitea_set_issue_labels`
- `gitea_snapshot_instance_fleet`
- `gitea_submit_pr_review`
- `gitea_update_pr_branch_by_merge`
- `gitea_validate_review_final_report`
+86 -15
View File
@@ -1193,6 +1193,31 @@ RECOGNIZED_GITEA_ENV_KEYS = frozenset({
"GITEA_IRRECOVERABLE_HMAC_SECRET",
"GITEA_FORCE_MCP_RUNTIME_CHECK",
"GITEA_FORCE_CLIENT_MANAGED",
# #975: the client-identity inputs the server actually consumes at startup
# (CLIENT_NAME_ENV / CLIENT_INSTANCE_ENV / CLIENT_SESSION_ENV in
# gitea_mcp_server). Production read them while this allowlist omitted them,
# so the peer-env scan classified them as unsupported overrides and the
# capability resolver refused every mutation fleet-wide. Named individually
# on purpose: no prefix is added, so an unrecognised GITEA_* override is
# still refused exactly as it was before.
"GITEA_MCP_CLIENT",
"GITEA_MCP_CLIENT_INSTANCE",
"GITEA_MCP_CLIENT_SESSION",
# #978: operator-approved fleet enrollment identity shared by one launch.
"GITEA_MCP_FLEET_RUN_ID",
"GITEA_MCP_PROCESS_IDENTITY",
# #978 B1/B2: launcher-sealed provenance + peer-visible worker identity so
# the instance-aware mutation gate can distinguish two legitimate
# application instances that share a profile from a true duplicate worker.
"GITEA_MCP_INSTANCE_PROVENANCE",
"GITEA_MCP_WORKER_IDENTITY",
"GITEA_MCP_GENERATION_ID",
# #975 review 652 B1: production also consumes HEARTBEAT_INTERVAL_ENV from
# mcp_worker_identity via gitea_mcp_server._start_worker_heartbeat. Omitting
# it reproduced the same unsupported-env → runtime_reconnect_required
# failure mode for the documented operator override. Named individually;
# no GITEA_* / GITEA_WORKER_* prefix is added.
"GITEA_WORKER_HEARTBEAT_INTERVAL_SECONDS",
})
RECOGNIZED_GITEA_ENV_PREFIXES = (
@@ -1220,24 +1245,70 @@ def get_unconsumed_gitea_env_overrides(env=None) -> dict[str, str]:
return unconsumed
def launcher_entry(profile_name, config_path=None):
def launcher_entry(
profile_name,
config_path=None,
*,
client_type=None,
client_instance_id=None,
fleet_run_id=None,
session_id=None,
launch_nonce=None,
):
"""Return a thin MCP launcher entry for *profile_name*.
Contains command/args and the GITEA_MCP_* / GITEA_CLIENT_MANAGED env vars — never a token
or password. Suitable for Claude / Gemini / Codex ``mcpServers`` blocks.
Contains command/args and the GITEA_MCP_* / GITEA_CLIENT_MANAGED env vars —
never a token or password. Suitable for Claude / Gemini / Codex
``mcpServers`` blocks.
#978 B1: production launches mint (or reuse) a trusted
``GITEA_MCP_CLIENT_INSTANCE`` so workers never fall back to the legacy
placeholder identity on the real serve path. Pass *client_instance_id* to
resume the same launch; omit it for a fresh launch (new ID).
"""
command, args = server_command()
return {
"gitea-tools": {
"command": command,
"args": args,
"env": {
"GITEA_MCP_CONFIG": config_path or DEFAULT_CONFIG_PATH,
"GITEA_MCP_PROFILE": profile_name,
"GITEA_CLIENT_MANAGED": "1",
},
}
}
import mcp_application_launcher as app_launcher
built = app_launcher.launcher_entry_for_profile(
profile_name,
client_type=client_type or "unknown",
config_path=config_path or DEFAULT_CONFIG_PATH,
client_instance_id=client_instance_id,
fleet_run_id=fleet_run_id,
session_id=session_id,
launch_nonce=launch_nonce,
server_key="gitea-tools",
)
# Public shape stays {server_key: {command, args, env}} for existing callers.
return {"gitea-tools": built["gitea-tools"]}
def multi_namespace_launcher_entries(
profile_by_namespace,
*,
client_type,
config_path=None,
client_instance_id=None,
fleet_run_id=None,
session_id=None,
launch_nonce=None,
):
"""Build production mcpServers for all five namespaces of one application launch.
One shared trusted ``client_instance_id`` is minted (or reused) and
propagated to every namespace worker env. See
:func:`mcp_application_launcher.build_application_mcp_servers`.
"""
import mcp_application_launcher as app_launcher
return app_launcher.build_application_mcp_servers(
profile_by_namespace,
client_type=client_type,
config_path=config_path or DEFAULT_CONFIG_PATH,
client_instance_id=client_instance_id,
fleet_run_id=fleet_run_id,
session_id=session_id,
launch_nonce=launch_nonce,
)
+451 -32
View File
@@ -1916,8 +1916,16 @@ def _verify_role_mutation_workspace(
if runtime_reasons:
raise RuntimeError("; ".join(runtime_reasons))
except Exception as exc:
if "stale-runtime:" in str(exc):
raise RuntimeError(str(exc))
# #975: ``_check_mcp_runtimes_diagnostics`` raises every reason it
# produces through this one RuntimeError, but only ``stale-runtime:``
# was re-raised here — an ``unsupported-env:`` reason was swallowed
# while still failing the capability resolver, so this preflight and
# the resolver disagreed about the identical diagnostic. Both
# authoritative prefixes now propagate the same way. This can only ever
# widen what is refused, never widen what is permitted.
message = str(exc)
if any(prefix in message for prefix in RUNTIME_DIAGNOSTIC_HARD_PREFIXES):
raise RuntimeError(message)
pass
role = _effective_workspace_role()
@@ -15714,6 +15722,10 @@ _WORKER_REGISTRY = None
_WORKER_IDENTITY: str | None = None
_WORKER_GENERATION: str | None = None
_WORKER_REGISTRATION_ATTEMPTED = False
#: #975: the one heartbeat supervisor for this process's registration. One
#: process registers exactly one worker identity, so there is exactly one
#: supervisor and it is never replaced.
_WORKER_HEARTBEAT_SUPERVISOR = None
#: Env a client launcher may set to name itself and its session. Absent values
#: are reported as unknown; they are never guessed at, because guessing is what
@@ -15721,6 +15733,10 @@ _WORKER_REGISTRATION_ATTEMPTED = False
CLIENT_NAME_ENV = "GITEA_MCP_CLIENT"
CLIENT_INSTANCE_ENV = "GITEA_MCP_CLIENT_INSTANCE"
CLIENT_SESSION_ENV = "GITEA_MCP_CLIENT_SESSION"
FLEET_RUN_ENV = "GITEA_MCP_FLEET_RUN_ID"
INSTANCE_PROVENANCE_ENV = "GITEA_MCP_INSTANCE_PROVENANCE"
WORKER_IDENTITY_ENV = "GITEA_MCP_WORKER_IDENTITY"
GENERATION_ID_ENV = "GITEA_MCP_GENERATION_ID"
def _worker_registry():
@@ -15744,13 +15760,57 @@ def _worker_registry():
def _client_identity_hints() -> dict:
"""What the launcher told us about itself. Unset fields stay unset."""
"""What the launcher told us about itself (#948 / #975 / #978).
``client_instance_id`` is established by the trusted application launcher
once per launch and shared by all five namespace workers. Workers inherit
that value from the production launcher env and never mint a trusted ID
themselves. When the launcher key is absent or untrusted we still register
under an explicit *legacy* placeholder so diagnostic reads work, but the
registration is marked incomplete and cannot authorize multi-instance fleet
mutation safety.
"""
import mcp_application_launcher as app_launcher
import mcp_fleet_snapshot
inherited = app_launcher.inherit_or_refuse_client_instance(os.environ)
client_name = (os.environ.get(CLIENT_NAME_ENV) or "").strip() or None
instance = {
"client_instance_id": inherited.get("client_instance_id"),
"provenance": inherited.get("provenance"),
"trusted": bool(inherited.get("trusted")),
"complete": bool(inherited.get("complete")),
"launcher_sealed": bool(inherited.get("launcher_sealed")),
}
# Legacy placeholder only when the launcher omitted a usable key — never
# invent a trusted ID from PID proximity or ordinary user env.
if not instance["client_instance_id"] or not instance["trusted"]:
if not instance["client_instance_id"]:
placeholder = f"legacy-pid-{os.getpid()}"
fallback = mcp_fleet_snapshot.assess_instance_identity(
placeholder, source="pid_fallback"
)
instance = {
"client_instance_id": fallback["client_instance_id"],
"provenance": fallback["provenance"],
"trusted": False,
"complete": False,
"launcher_sealed": False,
}
else:
# Malformed / untrusted user-supplied value: keep it for diagnosis
# but never mark trusted.
instance["trusted"] = False
instance["launcher_sealed"] = False
return {
"client_name": (os.environ.get(CLIENT_NAME_ENV) or "").strip() or None,
"client_instance_id": (os.environ.get(CLIENT_INSTANCE_ENV) or "").strip()
or f"pid-{os.getpid()}",
"client_name": client_name,
"client_instance_id": instance["client_instance_id"],
"instance_id_provenance": instance["provenance"],
"instance_identity_trusted": instance["trusted"],
"instance_launcher_sealed": instance.get("launcher_sealed", False),
"session_id": (os.environ.get(CLIENT_SESSION_ENV) or "").strip()
or f"proc-{os.getpid()}-{_process_boot_head_sha or 'nohead'}",
"fleet_run_id": (os.environ.get(FLEET_RUN_ENV) or "").strip() or None,
}
@@ -15782,24 +15842,71 @@ def _active_worker_identity() -> str | None:
identity = mcp_worker_identity.generate_worker_identity(
hints["client_name"], hints["session_id"]
)
role = _active_role_kind_safe()
profile_name = (
(os.environ.get(gitea_config.ENV_PROFILE) or "").strip() or None
)
# Namespace is the role namespace (author/reviewer/…) when known.
namespace = role if role in {
"author", "reviewer", "merger", "controller", "reconciler"
} else None
parity = None
try:
parity = _current_master_parity()
except Exception:
parity = None
outcome = registry.register(
worker_identity=identity,
client_name=hints["client_name"],
client_instance_id=hints["client_instance_id"],
session_id=hints["session_id"],
generation_id=generation,
role=_active_role_kind_safe(),
profile=(os.environ.get(gitea_config.ENV_PROFILE) or "").strip() or None,
role=role,
profile=profile_name,
namespace=namespace,
remote=(os.environ.get("GITEA_MCP_REMOTE") or "").strip() or None,
repository_binding=PROJECT_ROOT,
pid=os.getpid(),
transport=native.get("bound_transport"),
token_fingerprint=native.get("token_fingerprint"),
pid_alive_probe=issue_lock_store.is_process_alive,
fleet_run_id=hints.get("fleet_run_id"),
process_identity=f"pid-{os.getpid()}",
startup_revision=(parity or {}).get("startup_head")
or _process_boot_head_sha,
loaded_revision=(parity or {}).get("local_head")
or (parity or {}).get("current_head"),
parity_revision=(parity or {}).get("current_head"),
live_revision=(parity or {}).get("live_remote_head"),
instance_id_provenance=hints.get("instance_id_provenance"),
)
if outcome.get("registered"):
_WORKER_IDENTITY = identity
_WORKER_GENERATION = generation
# #978 B2: export worker/generation/instance into this process
# env so peer process scans (and the instance-aware mutation
# gate) can attribute live workers without treating shared
# profiles as duplicates. Never overwrite a launcher-sealed
# client_instance_id with a different value.
os.environ[WORKER_IDENTITY_ENV] = identity
os.environ[GENERATION_ID_ENV] = generation
if hints.get("client_instance_id") and not (
os.environ.get(CLIENT_INSTANCE_ENV) or ""
).strip():
os.environ[CLIENT_INSTANCE_ENV] = str(
hints["client_instance_id"]
)
# #975: registration is the only moment identity and fencing
# epoch are both known, so the heartbeat supervisor is started
# here. Without it ``last_heartbeat_at`` never left
# ``started_at`` and every healthy client lost ownership at the
# TTL.
_start_worker_heartbeat(
identity=identity,
fencing_epoch=outcome.get("fencing_epoch"),
hints=hints,
generation=generation,
)
return identity
if not outcome.get("collision"):
return None
@@ -15810,6 +15917,78 @@ def _active_worker_identity() -> str | None:
return None
def _start_worker_heartbeat(
*,
identity: str,
fencing_epoch,
hints: dict,
generation: str | None,
) -> None:
"""Attach a heartbeat supervisor to the registration just created (#975).
Never raises: a supervisor that cannot start leaves the row exactly as
``register()`` wrote it, which is the pre-#975 behaviour, rather than
failing the tool call that happened to trigger lazy registration.
Under pytest the thread is deliberately not started. Tests drive
``WorkerHeartbeatSupervisor`` directly with an injected clock, so the suite
proves the lifecycle without leaving background sqlite writers behind.
"""
global _WORKER_HEARTBEAT_SUPERVISOR
if _WORKER_HEARTBEAT_SUPERVISOR is not None:
return
registry = _worker_registry()
if registry is None or fencing_epoch is None:
return
try:
supervisor = mcp_worker_identity.WorkerHeartbeatSupervisor(
registry,
worker_identity=identity,
fencing_epoch=int(fencing_epoch),
session_id=hints.get("session_id"),
generation_id=generation,
client_name=hints.get("client_name"),
pid=os.getpid(),
ttl_seconds=mcp_worker_identity.DEFAULT_HEARTBEAT_TTL_SECONDS,
interval_seconds=os.environ.get(
mcp_worker_identity.HEARTBEAT_INTERVAL_ENV
),
)
_WORKER_HEARTBEAT_SUPERVISOR = supervisor
if not mcp_daemon_guard.is_pytest_runtime():
supervisor.start()
except Exception:
return
def _worker_heartbeat_status() -> dict:
"""Read-only heartbeat observability for ``gitea_get_runtime_context`` (#975).
Reports the unsupervised case explicitly rather than omitting the key, so an
operator can tell "no supervisor" apart from "supervisor with no beats yet".
"""
supervisor = _WORKER_HEARTBEAT_SUPERVISOR
if supervisor is None:
return {
"supervised": False,
"running": False,
"reasons": [
"no worker heartbeat supervisor is attached to this process; the "
"registration is not being renewed and will go stale at its TTL"
],
}
try:
status = dict(supervisor.status())
except Exception as exc:
return {
"supervised": True,
"running": False,
"reasons": [f"heartbeat status unavailable: {type(exc).__name__}: {exc}"],
}
status.setdefault("reasons", [])
return status
def _active_role_kind_safe() -> str | None:
"""Best-effort role for the registry record; never raises into a tool call.
@@ -19278,6 +19457,172 @@ def gitea_validate_review_final_report(
)
@mcp.tool()
def gitea_snapshot_instance_fleet(
remote: str = "dadeschools",
host: str | None = None,
org: str | None = None,
repo: str | None = None,
expected_manifest: list | dict | str | None = None,
include_historical: bool = True,
canonical_repository: str | None = None,
expected_live_revision: str | None = None,
) -> dict:
"""Read-only: authoritative instance-level fleet identity and health snapshot (#978).
Controller and reconciler only. Returns a point-in-time snapshot of every
registered namespace worker with instance attribution, heartbeat freshness,
and structured classification (expected/missing/unmanifested, duplicate
namespace workers, identity collisions, foreign repository, old revision,
historical dead rows). Sharing only a client type is never a duplicate.
Does not grant any mutation capability. Does not scan process tables or
open foreign databases for production evidence the worker registry is
the sole authority.
Args:
remote: Known instance 'dadeschools' or 'prgs'.
host: Optional host override.
org: Optional org override (audit context only).
repo: Optional repo override (audit context only).
expected_manifest: Optional approved fleet enrollment list
(list of dicts with client_instance_id / client_type / fleet_run_id).
When supplied, missing and unmanifested instances are classified.
include_historical: Include released/superseded rows (default True);
historical rows never make the live fleet unsafe by themselves.
canonical_repository: Expected repository binding path/slug for
foreign-repository classification.
expected_live_revision: When set, workers whose recorded revisions
differ are classified as old-revision.
"""
import json as _json
import mcp_fleet_snapshot
read_block = _profile_operation_gate("gitea.read")
if read_block:
return {
"success": False,
"read_only": True,
"reasons": read_block,
"permission_report": _permission_block_report("gitea.read"),
}
profile = get_profile()
role = _profile_role_kind(profile)
if role not in {"controller", "reconciler"}:
return {
"success": False,
"read_only": True,
"allowed": False,
"denied_role": role,
"required_roles": ["controller", "reconciler"],
"mutation_performed": False,
"reasons": [
f"gitea_snapshot_instance_fleet is restricted to controller and "
f"reconciler roles; active role_kind is {role!r}. Author, "
"reviewer, and merger profiles keep gitea.read for diagnosis "
"elsewhere but do not receive this fleet surface, and no "
"unrelated mutation permission is granted."
],
"exact_next_action": (
"Re-run from a prgs-controller or prgs-reconciler namespace."
),
}
manifest: list | None = None
if expected_manifest is not None:
raw = expected_manifest
if isinstance(raw, str):
try:
raw = _json.loads(raw)
except Exception as exc:
return {
"success": False,
"read_only": True,
"reasons": [
f"expected_manifest is not valid JSON: {type(exc).__name__}"
],
}
if isinstance(raw, dict):
# Accept {"instances": [...]} or a single instance object.
if "instances" in raw and isinstance(raw["instances"], list):
raw = raw["instances"]
else:
raw = [raw]
if not isinstance(raw, list):
return {
"success": False,
"read_only": True,
"reasons": ["expected_manifest must be a list of instance objects"],
}
manifest = raw
registry = _worker_registry()
if registry is None:
return {
"success": False,
"read_only": True,
"reasons": [
"worker registry is unavailable; cannot produce an authoritative "
"fleet snapshot (fail closed)"
],
"exact_next_action": (
"Ensure GITEA_WORKER_REGISTRY_DB is writable and re-run after "
"workers have registered."
),
}
try:
if include_historical:
records = registry.list_workers(status=None)
else:
records = registry.list_workers(status=mcp_worker_identity.STATUS_ACTIVE)
except Exception as exc:
return {
"success": False,
"read_only": True,
"reasons": [
f"failed to read worker registry: {type(exc).__name__}: {_redact(str(exc))}"
],
}
parity = None
try:
parity = _current_master_parity()
except Exception:
parity = None
live_rev = expected_live_revision or (parity or {}).get("live_remote_head")
canon = canonical_repository or PROJECT_ROOT
snapshot = mcp_fleet_snapshot.snapshot_instance_fleet(
records,
expected_manifest=manifest,
pid_alive_probe=issue_lock_store.is_process_alive,
canonical_repository=canon,
expected_live_revision=live_rev,
)
snapshot["role_kind"] = role
snapshot["profile"] = profile.get("profile_name")
snapshot["remote"] = _effective_remote(remote)
snapshot["repository"] = {
"org": org,
"repo": repo,
"canonical_repository": canon,
}
snapshot["mutation_performed"] = False
snapshot["permission_scope"] = {
"read_only": True,
"granted_operations": ["gitea.read"],
"denied_unrelated_mutations": True,
"note": (
"This capability is strictly observational. It does not authorize "
"branch, issue, PR, review, merge, or restart mutations."
),
}
return snapshot
@mcp.tool()
def gitea_get_runtime_context(
remote: str = "dadeschools",
@@ -19451,6 +19796,11 @@ def gitea_get_runtime_context(
"fencing_epoch": provenance_assessment["fencing_epoch"],
"conflicting_live_sessions": provenance_assessment["conflicting_live_sessions"],
"provenance_assessment": provenance_assessment,
# #975: whether this registration is actually being renewed. Read-only,
# and it grants nothing — ownership still comes from the attachment
# record above. It exists so "my heartbeat stopped" is diagnosable
# before the TTL turns it into session_attachment_missing.
"worker_heartbeat": _worker_heartbeat_status(),
"unconsumed_gitea_env": unconsumed_env,
"preflight_ready": preflight["preflight_ready"],
"preflight_block_reasons": preflight["preflight_block_reasons"],
@@ -21959,6 +22309,17 @@ def gitea_route_task_session(
# self-recovery was removed from the read-only path (was _trigger_mcp_auto_restart).
# Recovery is owned exclusively by the IDE/client reconnect path.
#: Every reason prefix ``_check_mcp_runtimes_diagnostics`` can emit. Callers
#: raise its reasons as one RuntimeError, and #975 found that the preflight
#: re-raise recognised only ``stale-runtime:``, silently dropping
#: ``unsupported-env:`` while the capability resolver still failed on it. This
#: lives beside the producer so a newly added reason family cannot be forgotten
#: by a distant re-raise predicate again.
RUNTIME_DIAGNOSTIC_HARD_PREFIXES: tuple[str, ...] = (
"stale-runtime:",
"unsupported-env:",
)
def _check_mcp_runtimes_diagnostics(task: str, matching_profiles: list[str]) -> list[str]:
"""Read-only: report missing or stale MCP runtimes (no config or process mutation).
@@ -22053,6 +22414,20 @@ def _check_mcp_runtimes_diagnostics(task: str, matching_profiles: list[str]) ->
)
peer_worker_identity = peer_env.get("GITEA_MCP_WORKER_IDENTITY")
peer_generation = peer_env.get("GITEA_MCP_GENERATION_ID")
# #978 B2: instance attribution for the live fleet gate. Two legitimate
# application instances may share a profile when each has a distinct
# trusted client_instance_id; only same-(instance, profile/namespace)
# duplicates remain blocked.
import mcp_fleet_snapshot as _fleet
peer_instance_raw = peer_env.get("GITEA_MCP_CLIENT_INSTANCE")
peer_instance_assess = _fleet.assess_instance_identity(peer_instance_raw)
peer_client_instance = (
peer_instance_assess["client_instance_id"]
if peer_instance_assess.get("trusted")
else None
)
peer_instance_trusted = bool(peer_instance_assess.get("trusted"))
for env_match in re.finditer(r'\b(GITEA_[A-Z0-9_]+)=([^\s]+)', env_out):
k, v = env_match.group(1), env_match.group(2)
@@ -22070,6 +22445,9 @@ def _check_mcp_runtimes_diagnostics(task: str, matching_profiles: list[str]) ->
"is_client_managed": is_client_managed,
"worker_identity": peer_worker_identity,
"generation_id": peer_generation,
"client_instance_id": peer_client_instance,
"instance_identity_trusted": peer_instance_trusted,
"raw_client_instance_id": peer_instance_raw,
}
if profile not in all_profile_procs:
all_profile_procs[profile] = []
@@ -22077,44 +22455,85 @@ def _check_mcp_runtimes_diagnostics(task: str, matching_profiles: list[str]) ->
running_profiles = {}
for profile, procs in all_profile_procs.items():
# #948 AC40/AC43: sharing a profile is legitimate — profile is a
# reusable capability definition, not a worker identity. What is *not*
# legitimate is reusing one worker identity, or two live sessions
# claiming one generation. Distinctness has to be proven, though:
# processes carrying no identity evidence are indistinguishable, so
# they stay classified as duplicates and keep the #686 wall intact.
identified = [p for p in procs if p.get("worker_identity")]
# #978 B2 / #948 AC40: sharing a profile is legitimate when each live
# process belongs to a distinct trusted application instance. The
# permanent model is instance-aware — not exactly_one_per_profile.
# Fail closed only when:
# * more than one live worker claims the same (client_instance_id,
# profile/namespace), or
# * worker identity / generation is reused, or
# * multiple processes share a profile without trusted instance
# evidence (indistinguishable → treat as duplicate).
# Group by trusted client_instance_id (untrusted/missing → one bucket).
by_instance: dict[str, list[dict]] = {}
for p in procs:
if p.get("instance_identity_trusted") and p.get("client_instance_id"):
key = str(p["client_instance_id"])
else:
key = "__untrusted_or_missing__"
by_instance.setdefault(key, []).append(p)
for instance_key, group in by_instance.items():
if len(group) <= 1:
continue
identified = [p for p in group if p.get("worker_identity")]
distinct_identities = {p["worker_identity"] for p in identified}
all_identified = len(identified) == len(procs)
all_identified = len(identified) == len(group)
duplicate_identity = len(identified) != len(distinct_identities)
contested_generation = any(
len({p["worker_identity"] for p in identified if p.get("generation_id") == gen}) > 1
for gen in {p.get("generation_id") for p in identified if p.get("generation_id")}
len(
{
p["worker_identity"]
for p in identified
if p.get("generation_id") == gen
}
)
if len(procs) > 1 and (
> 1
for gen in {
p.get("generation_id")
for p in identified
if p.get("generation_id")
}
)
# Same trusted (instance, profile) with >1 live process is always a
# duplicate-namespace-worker, even if worker identities differ —
# one application launch gets exactly one worker per namespace.
same_trusted_instance = instance_key != "__untrusted_or_missing__"
if same_trusted_instance or (
not all_identified or duplicate_identity or contested_generation
):
pids_str = ", ".join(str(p["pid"]) for p in procs)
if duplicate_identity or contested_generation:
pids_str = ", ".join(str(p["pid"]) for p in group)
if same_trusted_instance:
detail = (
"The same worker identity or generation is claimed more than once, "
"so these are genuine duplicates rather than independent workers."
f"More than one live worker for profile '{profile}' under "
f"client_instance_id {instance_key!r} (duplicate namespace "
"worker). Two legitimate application instances with "
"distinct trusted client_instance_id values may share a "
"profile; this collision does not."
)
elif duplicate_identity or contested_generation:
detail = (
"The same worker identity or generation is claimed more "
"than once, so these are genuine duplicates rather than "
"independent workers."
)
else:
detail = (
"Manual or duplicate launches defeat staleness detection and cannot "
"receive client stdio."
"Multiple processes share this profile without distinct "
"trusted client_instance_id evidence, so they cannot be "
"told apart from a duplicate MCP process. Relaunch via "
"the production application launcher so each instance "
"receives a trusted GITEA_MCP_CLIENT_INSTANCE."
)
reasons.append(
f"stale-runtime: Duplicate MCP server process(es) detected for profile '{profile}' (PIDs: {pids_str}). "
f"stale-runtime: Duplicate MCP server process(es) detected "
f"for profile '{profile}' (PIDs: {pids_str}). "
+ detail
)
# Otherwise: several independently identified workers share one profile.
# No reason is appended, deliberately. Every reason this function
# returns is raised as a hard RuntimeError by its callers, so recording
# legitimate concurrency here as "informational" would block exactly the
# case #948 exists to permit.
# Multiple trusted instances sharing one profile intentionally produce
# no reason here. Callers raise every reason as a hard RuntimeError, so
# treating legitimate multi-instance concurrency as "informational"
# would re-introduce the profile-only wall #978 removes.
client_procs = [p for p in procs if p["is_client_managed"]]
if client_procs:
client_procs.sort(key=lambda p: p["start_time"], reverse=True)
+393
View File
@@ -0,0 +1,393 @@
"""Production application launcher for multi-namespace MCP fleets (#978 B1).
One real LLM application launch mints exactly one trusted
``client_instance_id`` and propagates it to every Gitea MCP namespace worker
started for that launch. Workers never invent a trusted instance identity from
PID proximity, timestamps, or ordinary untrusted environment values.
This module is the production serve-path authority for instance identity.
Tests and fixtures may call the same functions, but production registration
receives the identity from the env this launcher builds — not from a hand-set
test-only helper that bypasses it.
"""
from __future__ import annotations
import os
import secrets
from typing import Any, Mapping
import mcp_fleet_snapshot as fleet
import mcp_worker_identity as mwi
#: Canonical MCP server names for the five role namespaces.
NAMESPACE_SERVER_NAMES: tuple[str, ...] = (
"gitea-author",
"gitea-reviewer",
"gitea-merger",
"gitea-controller",
"gitea-reconciler",
)
#: Map MCP server name → role/namespace kind.
SERVER_TO_NAMESPACE: dict[str, str] = {
"gitea-author": "author",
"gitea-reviewer": "reviewer",
"gitea-merger": "merger",
"gitea-controller": "controller",
"gitea-reconciler": "reconciler",
}
SANCTIONED_NAMESPACES: tuple[str, ...] = (
"author",
"reviewer",
"merger",
"controller",
"reconciler",
)
# Env the launcher may set on every worker of one application launch.
CLIENT_NAME_ENV = "GITEA_MCP_CLIENT"
CLIENT_INSTANCE_ENV = fleet.CLIENT_INSTANCE_ENV
FLEET_RUN_ENV = fleet.FLEET_RUN_ENV
CLIENT_SESSION_ENV = "GITEA_MCP_CLIENT_SESSION"
CLIENT_MANAGED_ENV = "GITEA_CLIENT_MANAGED"
PROFILE_ENV = "GITEA_MCP_PROFILE"
CONFIG_ENV = "GITEA_MCP_CONFIG"
INSTANCE_PROVENANCE_ENV = "GITEA_MCP_INSTANCE_PROVENANCE"
WORKER_IDENTITY_ENV = "GITEA_MCP_WORKER_IDENTITY"
GENERATION_ID_ENV = "GITEA_MCP_GENERATION_ID"
# Marker the trusted launcher alone writes; ordinary user env without this
# marker is never classified as launcher-trusted provenance.
LAUNCHER_PROVENANCE_VALUE = fleet.INSTANCE_ID_PROVENANCE_TRUSTED
def mint_application_launch(
client_type: str | None,
*,
launch_nonce: str | None = None,
fleet_run_id: str | None = None,
session_id: str | None = None,
now=None,
) -> dict[str, Any]:
"""Mint one trusted application-instance identity for a production launch.
Called exactly once per real application launch. The returned
``client_instance_id`` is injected into every namespace worker environment
for that launch. A second call (separate launch) yields a different ID.
"""
client = mwi.normalize_client_name(client_type)
instance_id = fleet.generate_client_instance_id(
client, launch_nonce=launch_nonce, now=now
)
assessment = fleet.assess_instance_identity(instance_id)
if not assessment["trusted"]:
# generate_client_instance_id always produces a trusted format; fail
# closed if that invariant ever breaks rather than shipping untrusted.
raise RuntimeError(
f"launcher produced untrusted client_instance_id {instance_id!r}: "
f"{assessment.get('reasons')}"
)
session = (session_id or "").strip() or f"launch-{secrets.token_hex(12)}"
return {
"client_type": client,
"client_instance_id": instance_id,
"instance_id_provenance": LAUNCHER_PROVENANCE_VALUE,
"instance_identity_trusted": True,
"fleet_run_id": (fleet_run_id or "").strip() or None,
"session_id": session,
"namespaces": list(SANCTIONED_NAMESPACES),
"namespace_server_names": list(NAMESPACE_SERVER_NAMES),
}
def namespace_worker_env(
*,
profile_name: str,
client_type: str | None,
client_instance_id: str,
config_path: str | None = None,
fleet_run_id: str | None = None,
session_id: str | None = None,
extra_env: Mapping[str, str] | None = None,
) -> dict[str, str]:
"""Build the environment for one namespace worker of a trusted launch.
Never invents a client_instance_id. The caller must supply the launch-minted
identity so all five workers receive the same value.
"""
assessment = fleet.assess_instance_identity(client_instance_id)
if not assessment["trusted"]:
raise ValueError(
"namespace_worker_env refuses untrusted client_instance_id "
f"{client_instance_id!r}: {assessment.get('reasons')}"
)
client = mwi.normalize_client_name(client_type)
env: dict[str, str] = {
PROFILE_ENV: str(profile_name),
CLIENT_MANAGED_ENV: "1",
"GITEA_MCP_CLIENT": client,
CLIENT_INSTANCE_ENV: assessment["client_instance_id"],
INSTANCE_PROVENANCE_ENV: LAUNCHER_PROVENANCE_VALUE,
}
if config_path:
env[CONFIG_ENV] = str(config_path)
if fleet_run_id:
env[FLEET_RUN_ENV] = str(fleet_run_id)
if session_id:
env[CLIENT_SESSION_ENV] = str(session_id)
if extra_env:
# Never let untrusted callers override the trusted instance keys.
protected = {
CLIENT_INSTANCE_ENV,
INSTANCE_PROVENANCE_ENV,
"GITEA_MCP_CLIENT",
CLIENT_MANAGED_ENV,
}
for key, value in extra_env.items():
if key in protected:
continue
env[str(key)] = str(value)
return env
def build_application_mcp_servers(
profile_by_namespace: Mapping[str, str],
*,
client_type: str | None,
config_path: str | None = None,
client_instance_id: str | None = None,
fleet_run_id: str | None = None,
session_id: str | None = None,
launch_nonce: str | None = None,
command: str | None = None,
args: list[str] | None = None,
now=None,
) -> dict[str, Any]:
"""Build a production ``mcpServers`` map for one application launch.
Mints one ``client_instance_id`` when *client_instance_id* is omitted (fresh
launch). When the caller supplies a previously minted trusted ID (resume of
the same launch / config rewrite), that ID is reused so reconnect keeps
attribution. A full new application restart omits the ID and receives a
fresh mint.
Every namespace server entry receives the **same** instance ID. Separate
calls with no supplied ID receive distinct IDs.
"""
missing = [
ns for ns in SANCTIONED_NAMESPACES if not profile_by_namespace.get(ns)
]
if missing:
raise ValueError(
"build_application_mcp_servers requires a profile for every "
f"sanctioned namespace; missing: {missing}"
)
if client_instance_id is None:
launch = mint_application_launch(
client_type,
launch_nonce=launch_nonce,
fleet_run_id=fleet_run_id,
session_id=session_id,
now=now,
)
else:
assessment = fleet.assess_instance_identity(client_instance_id)
if not assessment["trusted"]:
raise ValueError(
"refusing to propagate untrusted client_instance_id "
f"{client_instance_id!r} into production launch envs: "
f"{assessment.get('reasons')}"
)
launch = {
"client_type": mwi.normalize_client_name(client_type),
"client_instance_id": assessment["client_instance_id"],
"instance_id_provenance": LAUNCHER_PROVENANCE_VALUE,
"instance_identity_trusted": True,
"fleet_run_id": (fleet_run_id or "").strip() or None,
"session_id": (session_id or "").strip()
or f"launch-{secrets.token_hex(12)}",
"namespaces": list(SANCTIONED_NAMESPACES),
"namespace_server_names": list(NAMESPACE_SERVER_NAMES),
}
# Resolve command/args from the production server entry when not provided.
if command is None or args is None:
import gitea_config
cmd, cmd_args = gitea_config.server_command()
command = command or cmd
args = args if args is not None else list(cmd_args)
servers: dict[str, Any] = {}
shared_id = launch["client_instance_id"]
for namespace in SANCTIONED_NAMESPACES:
server_name = f"gitea-{namespace}"
profile = profile_by_namespace[namespace]
env = namespace_worker_env(
profile_name=profile,
client_type=launch["client_type"],
client_instance_id=shared_id,
config_path=config_path,
fleet_run_id=launch.get("fleet_run_id"),
session_id=launch.get("session_id"),
)
servers[server_name] = {
"command": command,
"args": list(args),
"env": env,
}
return {
"mcpServers": servers,
"launch": launch,
"client_instance_id": shared_id,
"client_type": launch["client_type"],
"namespaces": list(SANCTIONED_NAMESPACES),
"shared_instance_id_across_namespaces": True,
"namespace_count": len(SANCTIONED_NAMESPACES),
}
def launcher_entry_for_profile(
profile_name: str,
*,
client_type: str | None = None,
config_path: str | None = None,
client_instance_id: str | None = None,
fleet_run_id: str | None = None,
session_id: str | None = None,
server_key: str = "gitea-tools",
launch_nonce: str | None = None,
now=None,
) -> dict[str, Any]:
"""Thin single-server production launcher entry with trusted instance ID.
Used when only one namespace is being configured. Still mints (or reuses)
a trusted ``client_instance_id`` so production never relies on the legacy
placeholder identity for normal launches.
"""
import gitea_config
if client_instance_id is None:
launch = mint_application_launch(
client_type or "unknown",
launch_nonce=launch_nonce,
fleet_run_id=fleet_run_id,
session_id=session_id,
now=now,
)
client_instance_id = launch["client_instance_id"]
client = launch["client_type"]
fleet_run = launch.get("fleet_run_id")
session = launch.get("session_id")
else:
assessment = fleet.assess_instance_identity(client_instance_id)
if not assessment["trusted"]:
raise ValueError(
f"untrusted client_instance_id {client_instance_id!r}"
)
client = mwi.normalize_client_name(client_type)
fleet_run = (fleet_run_id or "").strip() or None
session = (session_id or "").strip() or None
client_instance_id = assessment["client_instance_id"]
command, args = gitea_config.server_command()
env = namespace_worker_env(
profile_name=profile_name,
client_type=client,
client_instance_id=client_instance_id,
config_path=config_path or gitea_config.DEFAULT_CONFIG_PATH,
fleet_run_id=fleet_run,
session_id=session,
)
return {
server_key: {
"command": command,
"args": args,
"env": env,
},
"client_instance_id": client_instance_id,
"client_type": client,
}
def collect_instance_ids_from_mcp_servers(
mcp_servers: Mapping[str, Any],
) -> dict[str, Any]:
"""Inspect a production mcpServers map for shared instance attribution.
Returns the unique set of client_instance_id values across gitea-* servers
and whether all five namespaces share exactly one trusted ID.
"""
ids: list[str] = []
by_server: dict[str, str | None] = {}
for name in NAMESPACE_SERVER_NAMES:
entry = mcp_servers.get(name) or {}
env = entry.get("env") or {}
raw = (env.get(CLIENT_INSTANCE_ENV) or "").strip() or None
by_server[name] = raw
if raw:
ids.append(raw)
unique = sorted(set(ids))
trusted = [
i
for i in unique
if fleet.assess_instance_identity(i)["trusted"]
]
return {
"server_instance_ids": by_server,
"unique_instance_ids": unique,
"trusted_instance_ids": trusted,
"shared_single_trusted_id": (
len(unique) == 1
and len(trusted) == 1
and all(by_server.get(n) == unique[0] for n in NAMESPACE_SERVER_NAMES)
),
"namespace_server_count": sum(
1 for n in NAMESPACE_SERVER_NAMES if n in mcp_servers
),
}
def inherit_or_refuse_client_instance(
env: Mapping[str, str] | None = None,
) -> dict[str, Any]:
"""Resolve instance identity for a worker process at serve time.
Production workers inherit the launcher-issued ID. They never mint a
trusted ID themselves. Missing / legacy / malformed values fail soft into
an untrusted assessment so registration can still record a diagnostic row
without authorizing multi-instance fleet mutation.
"""
source = dict(env if env is not None else os.environ)
raw = (source.get(CLIENT_INSTANCE_ENV) or "").strip() or None
provenance_marker = (source.get(INSTANCE_PROVENANCE_ENV) or "").strip()
assessment = fleet.assess_instance_identity(raw)
# Ordinary user-supplied values without launcher provenance marker are
# still format-checked by assess_instance_identity. When the format is
# trusted but the launcher marker is absent, keep the ID but note that
# provenance is not launcher-sealed (operator hand-set or legacy config).
if assessment["trusted"] and provenance_marker != LAUNCHER_PROVENANCE_VALUE:
assessment = dict(assessment)
assessment["launcher_sealed"] = False
assessment["reasons"] = list(assessment.get("reasons") or []) + [
f"{INSTANCE_PROVENANCE_ENV} is not {LAUNCHER_PROVENANCE_VALUE!r}; "
"identity format is valid but not sealed by the production launcher"
]
else:
assessment = dict(assessment)
assessment["launcher_sealed"] = bool(
assessment["trusted"]
and provenance_marker == LAUNCHER_PROVENANCE_VALUE
)
assessment["fleet_run_id"] = (source.get(FLEET_RUN_ENV) or "").strip() or None
assessment["session_id"] = (
(source.get(CLIENT_SESSION_ENV) or "").strip() or None
)
assessment["client_type"] = mwi.normalize_client_name(
(source.get("GITEA_MCP_CLIENT") or "").strip() or None
)
return assessment
+869
View File
@@ -0,0 +1,869 @@
"""Instance-level fleet identity and health snapshots (#978).
#948 established per-worker ownership; #975 made heartbeats keep those rows
live. Neither surface could enumerate the fleet at *instance* granularity:
which application launch owns which five namespace workers, whether two
Codex launches are distinct, or whether a live collision is real rather than
a shared client type.
This module is pure. Callers supply registry rows (and optional enrichments);
nothing here opens SQLite, scans process tables, or mutates state. Production
evidence for the fleet gate is the snapshot returned by the sanctioned
controller/reconciler tool that wraps this assessor.
Identity hierarchy (highest → lowest):
* ``client_type`` — application family (``codex``, ``claude_code``, …)
* ``client_instance_id`` — one running application launch (trusted launcher)
* ``fleet_run_id`` — operator-approved enrollment / canary cohort
* ``worker_id`` / ``worker_identity`` — one namespace worker process
* ``namespace`` — author | reviewer | merger | controller | reconciler
Multiple simultaneous instances of the same ``client_type`` are first-class.
Sharing only a profile or client type is never a duplicate.
"""
from __future__ import annotations
import hashlib
import secrets
from collections import defaultdict
from datetime import datetime, timezone
from typing import Any, Callable, Iterable, Mapping
import mcp_worker_identity as mwi
# --- Classification labels ------------------------------------------------
CLASS_EXPECTED = "expected_enrolled"
CLASS_MISSING = "missing_expected"
CLASS_UNMANIFESTED = "unmanifested"
CLASS_DUPLICATE_NAMESPACE = "duplicate_namespace_worker"
CLASS_INSTANCE_ID_COLLISION = "instance_id_collision"
CLASS_WORKER_ID_COLLISION = "worker_identity_collision"
CLASS_SESSION_COLLISION = "session_identity_collision"
CLASS_GENERATION_COLLISION = "generation_identity_collision"
CLASS_PROCESS_COLLISION = "process_identity_collision"
CLASS_PID_COLLISION = "pid_collision"
CLASS_OWNERSHIP_COLLISION = "ownership_fencing_collision"
CLASS_ORPHANED = "orphaned_unowned"
CLASS_UNKNOWN_CLIENT = "unknown_client"
CLASS_FOREIGN_REPOSITORY = "foreign_repository"
CLASS_OLD_REVISION = "old_revision"
CLASS_STALE_WORKER = "stale_orphaned_worker"
CLASS_LEGACY_INCOMPLETE = "legacy_incomplete_identity"
CLASS_HISTORICAL = "historical_dead"
CLASS_HEALTHY = "healthy"
#: Active blockers that make the live fleet unsafe for mutation-gated work.
ACTIVE_BLOCKER_CLASSES = frozenset(
{
CLASS_MISSING,
CLASS_UNMANIFESTED,
CLASS_DUPLICATE_NAMESPACE,
CLASS_INSTANCE_ID_COLLISION,
CLASS_WORKER_ID_COLLISION,
CLASS_SESSION_COLLISION,
CLASS_GENERATION_COLLISION,
CLASS_PROCESS_COLLISION,
CLASS_PID_COLLISION,
CLASS_OWNERSHIP_COLLISION,
CLASS_ORPHANED,
CLASS_UNKNOWN_CLIENT,
CLASS_FOREIGN_REPOSITORY,
CLASS_OLD_REVISION,
CLASS_STALE_WORKER,
CLASS_LEGACY_INCOMPLETE,
}
)
SANCTIONED_NAMESPACES = frozenset(
{"author", "reviewer", "merger", "controller", "reconciler"}
)
INSTANCE_ID_PROVENANCE_TRUSTED = "trusted_launcher"
INSTANCE_ID_PROVENANCE_LEGACY = "legacy_incomplete"
INSTANCE_ID_PROVENANCE_MISSING = "missing"
CLIENT_INSTANCE_ENV = "GITEA_MCP_CLIENT_INSTANCE"
FLEET_RUN_ENV = "GITEA_MCP_FLEET_RUN_ID"
PROCESS_IDENTITY_ENV = "GITEA_MCP_PROCESS_IDENTITY"
INSTANCE_PROVENANCE_ENV = "GITEA_MCP_INSTANCE_PROVENANCE"
_LEGACY_INSTANCE_PREFIXES = ("pid-", "proc-", "legacy-")
#: Trusted launcher-issued IDs use the reserved ``inst-`` prefix
#: (see :func:`generate_client_instance_id`). Ordinary user-supplied strings
#: without that prefix never count as trusted attribution.
_TRUSTED_INSTANCE_PREFIX = "inst-"
def _utc_now() -> datetime:
return datetime.now(timezone.utc)
def _ts(value: datetime) -> str:
return value.astimezone(timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ")
def generate_client_instance_id(
client_type: str | None,
*,
launch_nonce: str | None = None,
now: datetime | None = None,
) -> str:
"""Mint a distinct instance ID for one application launch (#978).
The trusted launcher (or host that starts all five namespace workers)
generates this once per launch and injects it as ``GITEA_MCP_CLIENT_INSTANCE``
into every worker environment. Workers never invent their own instance ID
from PID proximity or timestamps.
"""
client = mwi.normalize_client_name(client_type)
stamp = (now or _utc_now()).astimezone(timezone.utc).strftime(
mwi.IDENTITY_TIMESTAMP_FORMAT
)
nonce = launch_nonce if launch_nonce is not None else secrets.token_hex(16)
digest = hashlib.sha256(
f"{client}\x1f{stamp}\x1f{nonce}".encode("utf-8")
).hexdigest()[:12]
return f"inst-{client}-{stamp}-{digest}"
def assess_instance_identity(
raw_instance_id: str | None,
*,
source: str | None = None,
) -> dict[str, Any]:
"""Classify whether a client_instance_id is trusted enough for mutation.
Trusted instance IDs are non-empty, match the launcher-minted ``inst-…``
format from :func:`generate_client_instance_id`, and are not pre-#978
PID/proc/legacy placeholders. Ordinary user-supplied or malformed values
fail closed as untrusted so they cannot spoof multi-instance attribution.
Incomplete identities remain visible for diagnosis but cannot authorize
unsafe mutation.
"""
text = (raw_instance_id or "").strip()
if not text:
return {
"client_instance_id": None,
"complete": False,
"trusted": False,
"provenance": INSTANCE_ID_PROVENANCE_MISSING,
"reasons": [
"client_instance_id is missing; the trusted launcher must set "
f"{CLIENT_INSTANCE_ENV} once per application launch"
],
}
lowered = text.lower()
if lowered.startswith(_LEGACY_INSTANCE_PREFIXES) or source == "pid_fallback":
return {
"client_instance_id": text,
"complete": False,
"trusted": False,
"provenance": INSTANCE_ID_PROVENANCE_LEGACY,
"reasons": [
f"client_instance_id {text!r} is a legacy PID/process fallback, "
"not a trusted launcher-issued instance identity"
],
}
if not text.startswith(_TRUSTED_INSTANCE_PREFIX) or len(text) <= len(
_TRUSTED_INSTANCE_PREFIX
):
return {
"client_instance_id": text,
"complete": False,
"trusted": False,
"provenance": INSTANCE_ID_PROVENANCE_LEGACY,
"reasons": [
f"client_instance_id {text!r} is malformed or user-supplied and "
f"does not use the trusted launcher prefix "
f"{_TRUSTED_INSTANCE_PREFIX!r}; refusing trusted attribution"
],
}
return {
"client_instance_id": text,
"complete": True,
"trusted": True,
"provenance": INSTANCE_ID_PROVENANCE_TRUSTED,
"reasons": [],
}
def resolve_client_instance_from_env(
env: Mapping[str, str] | None = None,
*,
pid: int | None = None,
) -> dict[str, Any]:
"""Resolve instance identity from launcher env without inventing one.
When the trusted key is absent, return incomplete evidence rather than a
silent ``pid-<n>`` identity. Callers that still need a non-empty registry
key may choose a legacy placeholder deliberately; they must not treat it as
trusted.
"""
source = dict(env or {})
raw = (source.get(CLIENT_INSTANCE_ENV) or "").strip()
assessment = assess_instance_identity(raw or None)
assessment["fleet_run_id"] = (source.get(FLEET_RUN_ENV) or "").strip() or None
assessment["process_identity"] = (
(source.get(PROCESS_IDENTITY_ENV) or "").strip()
or (f"pid-{pid}" if pid is not None else None)
)
return assessment
def _public_worker(record: Mapping[str, Any]) -> dict[str, Any]:
return {
"worker_id": record.get("worker_identity") or record.get("worker_id"),
"worker_identity": record.get("worker_identity") or record.get("worker_id"),
"client_type": mwi.normalize_client_name(
record.get("client_name") or record.get("client_type")
),
"client_instance_id": record.get("client_instance_id"),
"fleet_run_id": record.get("fleet_run_id"),
"namespace": record.get("namespace"),
"profile": record.get("profile"),
"declared_role": record.get("role") or record.get("declared_role"),
"authenticated_account": record.get("authenticated_account"),
"session_id": record.get("session_id"),
"generation_id": record.get("generation_id"),
"process_identity": record.get("process_identity")
or (
f"pid-{record['pid']}"
if record.get("pid") is not None
else None
),
"pid": record.get("pid"),
"repository_binding": record.get("repository_binding"),
"remote": record.get("remote"),
"startup_revision": record.get("startup_revision"),
"loaded_revision": record.get("loaded_revision"),
"parity_revision": record.get("parity_revision"),
"live_revision": record.get("live_revision"),
"runtime_provenance": record.get("runtime_provenance")
or record.get("transport"),
"transport": record.get("transport"),
"started_at": record.get("started_at"),
"last_heartbeat_at": record.get("last_heartbeat_at"),
"heartbeat_ttl_seconds": record.get("heartbeat_ttl_seconds"),
"fencing_epoch": record.get("fencing_epoch"),
"status": record.get("status"),
"instance_id_provenance": record.get("instance_id_provenance"),
}
def _liveness(
record: Mapping[str, Any],
*,
now: datetime | None,
pid_alive_probe: Callable[[int | None], bool | None] | None,
) -> dict[str, Any]:
pid_alive = None
if pid_alive_probe is not None and record.get("pid") is not None:
try:
pid_alive = pid_alive_probe(record.get("pid"))
except Exception:
pid_alive = None
return mwi.WorkerRegistry.is_live(record, now=now, pid_alive=pid_alive)
def _consistency_token(rows: Iterable[Mapping[str, Any]], snapshot_at: str) -> str:
material = [snapshot_at]
for row in sorted(
rows,
key=lambda r: (
str(r.get("worker_identity") or ""),
str(r.get("last_heartbeat_at") or ""),
str(r.get("fencing_epoch") or ""),
),
):
material.append(
"|".join(
[
str(row.get("worker_identity") or ""),
str(row.get("client_instance_id") or ""),
str(row.get("status") or ""),
str(row.get("last_heartbeat_at") or ""),
str(row.get("fencing_epoch") or ""),
str(row.get("generation_id") or ""),
]
)
)
digest = hashlib.sha256("\n".join(material).encode("utf-8")).hexdigest()[:16]
return f"fleetrev-{digest}"
def build_worker_snapshot_row(
record: Mapping[str, Any],
*,
now: datetime | None = None,
pid_alive_probe: Callable[[int | None], bool | None] | None = None,
canonical_repository: str | None = None,
expected_live_revision: str | None = None,
heartbeat_supervised: bool | None = None,
) -> dict[str, Any]:
"""One point-in-time worker row for the fleet snapshot."""
stamp = now or _utc_now()
base = _public_worker(record)
identity = assess_instance_identity(
base.get("client_instance_id"),
source=record.get("instance_id_source"),
)
liveness = _liveness(record, now=stamp, pid_alive_probe=pid_alive_probe)
is_historical = str(record.get("status") or "") != mwi.STATUS_ACTIVE
live = bool(liveness.get("live")) and not is_historical
repo = (base.get("repository_binding") or "").strip() or None
foreign_repo = bool(
canonical_repository
and repo
and repo.rstrip("/") != str(canonical_repository).rstrip("/")
)
old_revision = False
if expected_live_revision:
for key in ("startup_revision", "loaded_revision", "parity_revision", "live_revision"):
rev = (base.get(key) or "").strip()
if rev and rev != expected_live_revision:
old_revision = True
break
ownership_state = "historical" if is_historical else (
"live" if live else "stale"
)
if live and not identity["trusted"]:
ownership_state = "live_untrusted_identity"
if live and not base.get("session_id"):
ownership_state = "orphaned"
mutation_safe = bool(
live
and identity["trusted"]
and not foreign_repo
and not old_revision
and ownership_state == "live"
and base.get("client_type") != mwi.UNKNOWN_CLIENT
)
restart_required = bool(
old_revision
or (live and not liveness.get("heartbeat_fresh", True))
)
return {
**base,
"instance_identity": identity,
"client_instance_id": identity["client_instance_id"] or base.get("client_instance_id"),
"instance_id_provenance": identity["provenance"],
"instance_identity_trusted": identity["trusted"],
"live": live,
"historical": is_historical,
"liveness": liveness,
"heartbeat": {
"registered": bool(base.get("last_heartbeat_at")),
"supervised": heartbeat_supervised,
"age_seconds": liveness.get("heartbeat_age_seconds"),
"ttl_seconds": liveness.get("heartbeat_ttl_seconds"),
"fresh": liveness.get("heartbeat_fresh"),
"last_heartbeat_at": base.get("last_heartbeat_at"),
},
"fencing": {
"fencing_epoch": base.get("fencing_epoch"),
"generation_id": base.get("generation_id"),
},
"ownership_state": ownership_state,
"foreign_repository": foreign_repo,
"old_revision": old_revision,
"stale": not live and not is_historical,
"restart_required": restart_required,
"mutation_safe": mutation_safe,
"conflicting_live_sessions": [],
}
def _collision_groups(
live_rows: list[dict[str, Any]],
key_fn,
) -> dict[str, list[dict[str, Any]]]:
groups: dict[str, list[dict[str, Any]]] = defaultdict(list)
for row in live_rows:
key = key_fn(row)
if key is None or key == "" or key == "None":
continue
groups[str(key)].append(row)
return {k: v for k, v in groups.items() if len(v) > 1}
def snapshot_instance_fleet(
records: Iterable[Mapping[str, Any]],
*,
expected_manifest: list[Mapping[str, Any]] | None = None,
now: datetime | None = None,
pid_alive_probe: Callable[[int | None], bool | None] | None = None,
canonical_repository: str | None = None,
expected_live_revision: str | None = None,
registry_revision: str | None = None,
known_client_types: Iterable[str] | None = None,
) -> dict[str, Any]:
"""Authoritative point-in-time fleet snapshot with classification (#978).
Historical dead rows are reported separately and never automatically make
the live fleet unsafe.
"""
stamp = now or _utc_now()
snapshot_at = _ts(stamp)
known = {
mwi.normalize_client_name(c)
for c in (known_client_types or mwi.CLIENT_ALIASES.values())
}
known.discard(mwi.UNKNOWN_CLIENT)
all_rows: list[dict[str, Any]] = []
for record in records:
all_rows.append(
build_worker_snapshot_row(
record,
now=stamp,
pid_alive_probe=pid_alive_probe,
canonical_repository=canonical_repository,
expected_live_revision=expected_live_revision,
)
)
live_rows = [r for r in all_rows if r["live"]]
historical_rows = [r for r in all_rows if r["historical"]]
stale_rows = [r for r in all_rows if r["stale"]]
# --- identity collisions among live workers ---
findings: list[dict[str, Any]] = []
def _finding(
classification: str,
*,
severity: str,
workers: list[dict[str, Any]] | None = None,
instance_ids: list[str] | None = None,
detail: str,
active_blocker: bool,
) -> None:
findings.append(
{
"classification": classification,
"severity": severity,
"active_blocker": active_blocker,
"detail": detail,
"client_instance_ids": instance_ids or sorted(
{
str(w.get("client_instance_id"))
for w in (workers or [])
if w.get("client_instance_id")
}
),
"worker_identities": [
w.get("worker_identity") for w in (workers or [])
],
}
)
# Duplicate worker identity (should not happen with PK, still detect)
for wid, group in _collision_groups(
live_rows, lambda r: r.get("worker_identity")
).items():
_finding(
CLASS_WORKER_ID_COLLISION,
severity="blocker",
workers=group,
detail=f"worker identity {wid!r} is claimed by {len(group)} live workers",
active_blocker=True,
)
# Reused session identity across live workers
for sid, group in _collision_groups(live_rows, lambda r: r.get("session_id")).items():
# Same session may appear once; collision only when multiple workers share it
# across different worker identities (always true for group size > 1).
_finding(
CLASS_SESSION_COLLISION,
severity="blocker",
workers=group,
detail=f"session identity {sid!r} is reused by {len(group)} live workers",
active_blocker=True,
)
# Generation claimed by multiple live sessions/workers is a conflict when
# the workers are not the five sanctioned namespaces of one instance.
for gen, group in _collision_groups(
live_rows, lambda r: r.get("generation_id")
).items():
namespaces = {g.get("namespace") for g in group if g.get("namespace")}
instance_ids = {g.get("client_instance_id") for g in group}
# Multiple workers under one generation is only valid if they share one
# instance and distinct namespaces. Same generation + same namespace = bad.
by_ns: dict[str, list] = defaultdict(list)
for g in group:
by_ns[str(g.get("namespace") or "")].append(g)
ns_dups = {ns: rows for ns, rows in by_ns.items() if ns and len(rows) > 1}
if ns_dups or len(instance_ids) > 1:
_finding(
CLASS_GENERATION_COLLISION,
severity="blocker",
workers=group,
detail=(
f"generation {gen!r} is contested across namespaces/instances "
f"(namespaces={sorted(namespaces)}, "
f"instances={sorted(str(i) for i in instance_ids if i)})"
),
active_blocker=True,
)
# Process identity / PID collisions across distinct workers
for proc, group in _collision_groups(
live_rows, lambda r: r.get("process_identity")
).items():
if len({r.get("worker_identity") for r in group}) > 1:
_finding(
CLASS_PROCESS_COLLISION,
severity="blocker",
workers=group,
detail=f"process identity {proc!r} is shared by distinct live workers",
active_blocker=True,
)
for pid, group in _collision_groups(live_rows, lambda r: r.get("pid")).items():
if len({r.get("worker_identity") for r in group}) > 1:
_finding(
CLASS_PID_COLLISION,
severity="blocker",
workers=group,
detail=f"PID {pid} is shared by distinct live workers",
active_blocker=True,
)
# Fencing/ownership: same fencing epoch on different workers of different instances
for epoch, group in _collision_groups(
live_rows,
lambda r: (
f"{r.get('generation_id')}:{r.get('fencing_epoch')}"
if r.get("generation_id") is not None and r.get("fencing_epoch") is not None
else None
),
).items():
if len({r.get("client_instance_id") for r in group}) > 1:
_finding(
CLASS_OWNERSHIP_COLLISION,
severity="blocker",
workers=group,
detail=(
f"fencing token {epoch!r} spans more than one client_instance_id"
),
active_blocker=True,
)
# Per-instance grouping
by_instance: dict[str, list[dict[str, Any]]] = defaultdict(list)
unkeyed_live: list[dict[str, Any]] = []
for row in live_rows:
iid = row.get("client_instance_id")
if not iid:
unkeyed_live.append(row)
continue
by_instance[str(iid)].append(row)
instances: list[dict[str, Any]] = []
for iid, workers in sorted(by_instance.items()):
client_types = sorted({w.get("client_type") for w in workers if w.get("client_type")})
trusted = all(w.get("instance_identity_trusted") for w in workers)
namespaces = [w.get("namespace") for w in workers]
ns_counts: dict[str, int] = defaultdict(int)
for ns in namespaces:
if ns:
ns_counts[str(ns)] += 1
dup_ns = sorted(ns for ns, n in ns_counts.items() if n > 1)
if dup_ns:
_finding(
CLASS_DUPLICATE_NAMESPACE,
severity="blocker",
workers=[w for w in workers if w.get("namespace") in dup_ns],
instance_ids=[iid],
detail=(
f"instance {iid!r} has more than one live worker for "
f"namespace(s) {dup_ns}"
),
active_blocker=True,
)
# Live reuse of one instance ID with incompatible client types
if len(client_types) > 1:
_finding(
CLASS_INSTANCE_ID_COLLISION,
severity="blocker",
workers=workers,
instance_ids=[iid],
detail=(
f"client_instance_id {iid!r} is live under multiple client "
f"types {client_types}"
),
active_blocker=True,
)
if not trusted:
_finding(
CLASS_LEGACY_INCOMPLETE,
severity="blocker",
workers=workers,
instance_ids=[iid],
detail=(
f"instance {iid!r} lacks trusted launcher-issued instance "
"identity; diagnostic reads remain available"
),
active_blocker=True,
)
unknown = [w for w in workers if w.get("client_type") == mwi.UNKNOWN_CLIENT]
if unknown:
_finding(
CLASS_UNKNOWN_CLIENT,
severity="blocker",
workers=unknown,
instance_ids=[iid],
detail=f"instance {iid!r} has worker(s) with unknown client_type",
active_blocker=True,
)
foreign = [w for w in workers if w.get("foreign_repository")]
if foreign:
_finding(
CLASS_FOREIGN_REPOSITORY,
severity="blocker",
workers=foreign,
instance_ids=[iid],
detail=f"instance {iid!r} has foreign-repository workers",
active_blocker=True,
)
old = [w for w in workers if w.get("old_revision")]
if old:
_finding(
CLASS_OLD_REVISION,
severity="blocker",
workers=old,
instance_ids=[iid],
detail=f"instance {iid!r} has old-revision workers",
active_blocker=True,
)
orphans = [w for w in workers if w.get("ownership_state") == "orphaned"]
if orphans:
_finding(
CLASS_ORPHANED,
severity="blocker",
workers=orphans,
instance_ids=[iid],
detail=f"instance {iid!r} has orphaned/unowned workers",
active_blocker=True,
)
instances.append(
{
"client_instance_id": iid,
"client_types": client_types,
"client_type": client_types[0] if len(client_types) == 1 else None,
"fleet_run_ids": sorted(
{w.get("fleet_run_id") for w in workers if w.get("fleet_run_id")}
),
"worker_count": len(workers),
"namespaces": sorted({n for n in namespaces if n}),
"namespace_counts": dict(ns_counts),
"duplicate_namespaces": dup_ns,
"trusted_instance_identity": trusted,
"workers": workers,
"mutation_safe": all(w.get("mutation_safe") for w in workers)
and not dup_ns
and trusted,
}
)
for row in unkeyed_live:
_finding(
CLASS_LEGACY_INCOMPLETE,
severity="blocker",
workers=[row],
detail="live worker has no client_instance_id",
active_blocker=True,
)
if row.get("client_type") == mwi.UNKNOWN_CLIENT:
_finding(
CLASS_UNKNOWN_CLIENT,
severity="blocker",
workers=[row],
detail="live worker has unknown client_type and no instance id",
active_blocker=True,
)
for row in stale_rows:
_finding(
CLASS_STALE_WORKER,
severity="warning",
workers=[row],
detail=(
f"worker {row.get('worker_identity')!r} is active in the registry "
"but not live (stale heartbeat or dead pid)"
),
active_blocker=True,
)
for row in historical_rows:
_finding(
CLASS_HISTORICAL,
severity="info",
workers=[row],
detail=(
f"historical registration {row.get('worker_identity')!r} "
f"(status={row.get('status')!r}) is not an active blocker"
),
active_blocker=False,
)
# Manifest comparison
expected = list(expected_manifest or [])
expected_ids = {
str(item.get("client_instance_id")).strip()
for item in expected
if (item.get("client_instance_id") or "").strip()
}
live_ids = set(by_instance.keys())
missing_ids = sorted(expected_ids - live_ids)
unmanifested_ids = sorted(live_ids - expected_ids) if expected_ids else []
for iid in missing_ids:
_finding(
CLASS_MISSING,
severity="blocker",
instance_ids=[iid],
detail=f"expected enrolled instance {iid!r} is missing from the live fleet",
active_blocker=True,
)
for iid in unmanifested_ids:
_finding(
CLASS_UNMANIFESTED,
severity="blocker",
instance_ids=[iid],
workers=by_instance.get(iid, []),
detail=(
f"live instance {iid!r} is not on the approved fleet manifest "
"(unmanifested)"
),
active_blocker=True,
)
# Same client_type multi-instance is healthy when each has distinct instance IDs
by_type: dict[str, list[str]] = defaultdict(list)
for inst in instances:
for ct in inst.get("client_types") or []:
by_type[str(ct)].append(inst["client_instance_id"])
multi_instance_same_type = {
ct: ids for ct, ids in by_type.items() if len(ids) > 1
}
active_blockers = [f for f in findings if f.get("active_blocker")]
historical_only = [f for f in findings if f.get("classification") == CLASS_HISTORICAL]
live_safe = not active_blockers
consistency = registry_revision or _consistency_token(all_rows, snapshot_at)
return {
"success": True,
"read_only": True,
"snapshot_at": snapshot_at,
"consistency_token": consistency,
"registry_revision": consistency,
"live_worker_count": len(live_rows),
"historical_worker_count": len(historical_rows),
"stale_worker_count": len(stale_rows),
"instance_count": len(instances),
"workers": all_rows,
"live_workers": live_rows,
"historical_workers": historical_rows,
"stale_workers": stale_rows,
"instances": instances,
"multi_instance_same_client_type": multi_instance_same_type,
"same_client_type_not_duplicate": True,
"expected_manifest": [
{
"client_instance_id": item.get("client_instance_id"),
"client_type": item.get("client_type"),
"fleet_run_id": item.get("fleet_run_id"),
"namespaces": item.get("namespaces"),
}
for item in expected
],
"missing_expected_instance_ids": missing_ids,
"unmanifested_instance_ids": unmanifested_ids,
"findings": findings,
"active_blockers": active_blockers,
"historical_findings": historical_only,
"live_fleet_safe": live_safe,
"mutation_safe": live_safe and all(
inst.get("mutation_safe") for inst in instances
)
if instances
else live_safe,
"classification_model": {
"exactly_one_process_per_profile": False,
"exactly_one_instance_per_client_type": False,
"multiple_instances_per_client_type": True,
"duplicate_requires": [
"live client_instance_id collision",
"duplicate namespace worker within one instance",
"reused worker/session/generation/process/pid/fencing identity",
"cross-instance ownership collision",
"unmanifested instance when a manifest is required",
],
},
"reasons": [f["detail"] for f in active_blockers],
}
def compare_snapshot_heartbeats(
earlier: Mapping[str, Any],
later: Mapping[str, Any],
) -> dict[str, Any]:
"""Prove heartbeat continuity and stable ownership across two snapshots."""
earlier_live = {
w.get("worker_identity"): w for w in earlier.get("live_workers") or []
}
later_live = {
w.get("worker_identity"): w for w in later.get("live_workers") or []
}
shared = sorted(set(earlier_live) & set(later_live))
continuity: list[dict[str, Any]] = []
stable_ownership = True
for wid in shared:
a = earlier_live[wid]
b = later_live[wid]
same_instance = a.get("client_instance_id") == b.get("client_instance_id")
same_session = a.get("session_id") == b.get("session_id")
same_generation = a.get("generation_id") == b.get("generation_id")
hb_advanced_or_equal = True
if a.get("last_heartbeat_at") and b.get("last_heartbeat_at"):
hb_advanced_or_equal = b["last_heartbeat_at"] >= a["last_heartbeat_at"]
if not (same_instance and same_session and same_generation):
stable_ownership = False
continuity.append(
{
"worker_identity": wid,
"same_client_instance_id": same_instance,
"same_session_id": same_session,
"same_generation_id": same_generation,
"heartbeat_non_decreasing": hb_advanced_or_equal,
"earlier_heartbeat": a.get("last_heartbeat_at"),
"later_heartbeat": b.get("last_heartbeat_at"),
}
)
return {
"shared_live_workers": shared,
"continuity": continuity,
"stable_ownership": stable_ownership
and all(c["heartbeat_non_decreasing"] for c in continuity),
"dropped_workers": sorted(set(earlier_live) - set(later_live)),
"new_workers": sorted(set(later_live) - set(earlier_live)),
}
+483 -2
View File
@@ -29,6 +29,7 @@ implements:
from __future__ import annotations
import atexit
import hashlib
import os
import re
@@ -123,6 +124,103 @@ DEFAULT_REGISTRY_PATH = os.path.expanduser(
#: worker within a single operator coffee break.
DEFAULT_HEARTBEAT_TTL_SECONDS = 900.0
#: Optional operator override for the beat interval, in seconds. Clamped by
#: :func:`heartbeat_interval_for` so it can never be set at or above the TTL.
HEARTBEAT_INTERVAL_ENV = "GITEA_WORKER_HEARTBEAT_INTERVAL_SECONDS"
#: Never beat faster than this, so a misconfigured interval cannot turn the
#: supervisor into a busy sqlite writer.
MIN_HEARTBEAT_INTERVAL_SECONDS = 1.0
def heartbeat_interval_for(
ttl_seconds: float = DEFAULT_HEARTBEAT_TTL_SECONDS,
override: str | float | None = None,
) -> float:
"""Beat interval for *ttl_seconds*: one third of the TTL, capped at a half.
#975. A third means two consecutive beats can be lost before the row is
allowed to look stale, and the half-TTL cap is a hard ceiling so no
override can produce an interval that expires the registration it is
supposed to renew. That is the whole liveness contract: a healthy worker
stays owned, and a worker that has genuinely stopped beating still expires
on the configured TTL — this function never touches the TTL itself.
"""
try:
ttl = float(ttl_seconds)
except (TypeError, ValueError):
ttl = DEFAULT_HEARTBEAT_TTL_SECONDS
if not ttl > 0:
ttl = DEFAULT_HEARTBEAT_TTL_SECONDS
interval = ttl / 3.0
if override is not None:
candidate = str(override).strip()
if candidate:
try:
parsed = float(candidate)
except (TypeError, ValueError):
parsed = None
if parsed is not None and parsed > 0:
interval = parsed
ceiling = ttl / 2.0
if interval > ceiling:
interval = ceiling
if interval < MIN_HEARTBEAT_INTERVAL_SECONDS:
interval = min(MIN_HEARTBEAT_INTERVAL_SECONDS, ceiling)
return interval
#: Recorded-vs-presented pairs compared by :func:`_heartbeat_expectation_drift`.
_HEARTBEAT_EXPECTATION_FIELDS = (
"session_id",
"generation_id",
"client_name",
"pid",
)
def _heartbeat_expectation_drift(
record: dict[str, Any],
*,
expected_session_id: str | None = None,
expected_generation_id: str | None = None,
expected_client_name: str | None = None,
expected_pid: int | None = None,
) -> list[tuple[str, Any, Any]]:
"""Return ``(field, recorded, presented)`` for every mismatched expectation.
Only supplied expectations are compared, so a caller that presents nothing
gets the pre-#975 behaviour. ``client_name`` is compared through
:func:`normalize_client_name` so one application's namespaces stay one
client while genuinely different clients stay distinct.
"""
presented = {
"session_id": expected_session_id,
"generation_id": expected_generation_id,
"client_name": (
normalize_client_name(expected_client_name)
if expected_client_name is not None
else None
),
"pid": expected_pid,
}
drift: list[tuple[str, Any, Any]] = []
for field in _HEARTBEAT_EXPECTATION_FIELDS:
want = presented[field]
if want is None:
continue
got = record.get(field)
if field == "pid":
match = got is not None and int(got) == int(want)
else:
match = got == want
if not match:
drift.append((field, got, want))
return drift
STATUS_ACTIVE = "active"
STATUS_SUPERSEDED = "superseded"
STATUS_RELEASED = "released"
@@ -203,8 +301,22 @@ CREATE INDEX IF NOT EXISTS idx_worker_session
ON worker_registrations(session_id, status);
CREATE INDEX IF NOT EXISTS idx_worker_profile
ON worker_registrations(profile, status);
CREATE INDEX IF NOT EXISTS idx_worker_instance
ON worker_registrations(client_instance_id, status);
"""
#: #978 optional columns added without rewriting historical rows.
_SCHEMA_OPTIONAL_COLUMNS: tuple[tuple[str, str], ...] = (
("fleet_run_id", "TEXT"),
("authenticated_account", "TEXT"),
("process_identity", "TEXT"),
("startup_revision", "TEXT"),
("loaded_revision", "TEXT"),
("parity_revision", "TEXT"),
("live_revision", "TEXT"),
("instance_id_provenance", "TEXT"),
)
class WorkerRegistryError(RuntimeError):
"""Raised for registry misuse that is a programming error, not a refusal."""
@@ -536,6 +648,17 @@ class WorkerRegistry:
conn = self._connect()
try:
conn.executescript(_SCHEMA_SQL)
existing = {
row[1]
for row in conn.execute(
"PRAGMA table_info(worker_registrations)"
).fetchall()
}
for name, decl in _SCHEMA_OPTIONAL_COLUMNS:
if name not in existing:
conn.execute(
f"ALTER TABLE worker_registrations ADD COLUMN {name} {decl}"
)
conn.commit()
finally:
conn.close()
@@ -660,6 +783,14 @@ class WorkerRegistry:
heartbeat_ttl_seconds: float = DEFAULT_HEARTBEAT_TTL_SECONDS,
now: datetime | None = None,
pid_alive_probe=None,
fleet_run_id: str | None = None,
authenticated_account: str | None = None,
process_identity: str | None = None,
startup_revision: str | None = None,
loaded_revision: str | None = None,
parity_revision: str | None = None,
live_revision: str | None = None,
instance_id_provenance: str | None = None,
) -> dict[str, Any]:
"""Atomically register one worker identity.
@@ -668,6 +799,10 @@ class WorkerRegistry:
to mint a different identity and register that instead (AC32). The
existing registration is returned untouched so the caller can see what
it collided with.
#978 also fails closed when a *live* registration already holds the same
``client_instance_id`` for the same namespace: two workers of one
namespace cannot share one instance.
"""
parsed = parse_worker_identity(worker_identity)
if not parsed["valid"]:
@@ -685,6 +820,9 @@ class WorkerRegistry:
}
stamp = _ts(now or _utc_now())
proc_id = process_identity or (
f"pid-{int(pid)}" if pid is not None else None
)
with self._tx() as conn:
existing = conn.execute(
"SELECT * FROM worker_registrations WHERE worker_identity = ?",
@@ -721,6 +859,41 @@ class WorkerRegistry:
),
}
# #978: refuse a second live worker for the same (instance, namespace).
if namespace and client_instance_id:
peers = conn.execute(
"SELECT * FROM worker_registrations "
"WHERE client_instance_id = ? AND namespace = ? AND status = ?",
(client_instance_id, namespace, STATUS_ACTIVE),
).fetchall()
for peer_row in peers:
peer = self._row_to_record(peer_row)
pid_alive = (
pid_alive_probe(peer.get("pid"))
if pid_alive_probe is not None
else None
)
if self.is_live(peer, now=now, pid_alive=pid_alive)["live"]:
return {
"success": False,
"registered": False,
"mutation_performed": False,
"blocker_kind": BLOCKER_CONFLICTING_SESSIONS,
"collision": True,
"collision_kind": "duplicate_namespace_worker",
"existing_registration": _public_record(peer),
"reasons": [
f"client_instance_id {client_instance_id!r} already "
f"has a live worker for namespace {namespace!r} "
f"({peer.get('worker_identity')!r}); #978 fails closed "
"rather than registering a second worker"
],
"exact_next_action": (
"Stop the extra namespace worker, or use a distinct "
"client_instance_id for a separate application launch."
),
}
epoch = self._next_epoch(conn, generation_id)
conn.execute(
"""
@@ -729,8 +902,11 @@ class WorkerRegistry:
generation_id, role, profile, namespace, remote,
repository_binding, pid, transport, token_fingerprint,
started_at, last_heartbeat_at, heartbeat_ttl_seconds,
fencing_epoch, status
) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)
fencing_epoch, status,
fleet_run_id, authenticated_account, process_identity,
startup_revision, loaded_revision, parity_revision,
live_revision, instance_id_provenance
) VALUES (?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)
""",
(
worker_identity,
@@ -751,6 +927,14 @@ class WorkerRegistry:
float(heartbeat_ttl_seconds),
epoch,
STATUS_ACTIVE,
(fleet_run_id or "").strip() or None,
(authenticated_account or "").strip() or None,
proc_id,
(startup_revision or "").strip() or None,
(loaded_revision or "").strip() or None,
(parity_revision or "").strip() or None,
(live_revision or "").strip() or None,
(instance_id_provenance or "").strip() or None,
),
)
row = conn.execute(
@@ -786,11 +970,25 @@ class WorkerRegistry:
worker_identity: str,
fencing_epoch: int,
now: datetime | None = None,
expected_session_id: str | None = None,
expected_generation_id: str | None = None,
expected_client_name: str | None = None,
expected_pid: int | None = None,
) -> dict[str, Any]:
"""Renew only the owning registration (#948 AC11).
A stale epoch is refused rather than silently renewed, so a superseded
session that resumes cannot heartbeat its way back into ownership.
#975 adds optional keyword-only *expectations*. Every one supplied must
match the recorded row or the renewal is refused. They exist because
identity plus epoch cannot express "renew the row I registered, and only
that row": a recycled PID, or a second session of the same client, would
otherwise be renewable by the wrong beater. They are fencing tokens,
never assertions — supplying one can only cause a refusal, never grant
anything, and omitting them preserves the pre-#975 behaviour exactly.
Drift reports the existing ``BLOCKER_FENCED`` rather than a new
``blocker_kind``, because consumers switch on that value.
"""
stamp = _ts(now or _utc_now())
with self._tx() as conn:
@@ -807,6 +1005,30 @@ class WorkerRegistry:
"reasons": [f"no registration for {worker_identity!r}"],
}
record = self._row_to_record(row)
drift = _heartbeat_expectation_drift(
record,
expected_session_id=expected_session_id,
expected_generation_id=expected_generation_id,
expected_client_name=expected_client_name,
expected_pid=expected_pid,
)
if drift:
return {
"success": False,
"renewed": False,
"mutation_performed": False,
"blocker_kind": BLOCKER_FENCED,
"expectation_drift": drift,
"reasons": [
"presented worker expectations do not match the recorded "
"registration, so this beater does not own the row: "
+ "; ".join(
f"{field} recorded {recorded!r}, presented {presented!r}"
for field, recorded, presented in drift
)
+ " (#975)"
],
}
if record["status"] != STATUS_ACTIVE:
return {
"success": False,
@@ -1013,8 +1235,17 @@ def _public_record(record: dict[str, Any]) -> dict[str, Any]:
"transport": record.get("transport"),
"started_at": record.get("started_at"),
"last_heartbeat_at": record.get("last_heartbeat_at"),
"heartbeat_ttl_seconds": record.get("heartbeat_ttl_seconds"),
"fencing_epoch": record.get("fencing_epoch"),
"status": record.get("status"),
"fleet_run_id": record.get("fleet_run_id"),
"authenticated_account": record.get("authenticated_account"),
"process_identity": record.get("process_identity"),
"startup_revision": record.get("startup_revision"),
"loaded_revision": record.get("loaded_revision"),
"parity_revision": record.get("parity_revision"),
"live_revision": record.get("live_revision"),
"instance_id_provenance": record.get("instance_id_provenance"),
}
@@ -1430,3 +1661,253 @@ def resolve_bound_remote(
"explicitly to avoid host drift (#948)."
],
}
# --- Production heartbeat lifecycle (#975) --------------------------------
#
# #948 delivered ``WorkerRegistry.heartbeat()`` and nothing ever called it, so
# ``last_heartbeat_at`` stayed pinned to ``started_at`` for every production
# worker and ``heartbeat_ttl_seconds`` stopped being a liveness window at all —
# it became a hard cap on how long any client could stay attached. This
# supervisor is the missing caller.
#
# It is a daemon thread rather than an asyncio task or a per-request refresh
# because the renewal has to survive an *idle* session: a client waiting on a
# lease with no tool call in flight must not lose ownership, which rules out
# request-driven refresh. Registration itself happens lazily inside a tool call
# on the server's event loop, and the registry performs blocking
# ``BEGIN IMMEDIATE`` sqlite writes, which must not run on that loop.
# ``daemon=True`` is deliberate: a hard kill takes the thread down with the
# process, so a dead worker still goes stale on the normal TTL.
#: Refusals that mean this worker no longer owns its row. Beating again could
#: only ever be an attempt to renew ownership it has already lost, so the
#: supervisor stops permanently instead of retrying.
TERMINAL_HEARTBEAT_BLOCKERS = frozenset({BLOCKER_FENCED, BLOCKER_NO_ATTACHMENT})
class WorkerHeartbeatSupervisor:
"""Periodically renew exactly one worker registration.
Contract:
* It never registers. A supervisor exists only for an already-registered
identity, so it cannot create a second registration or a second identity
system.
* Every beat presents the full expectation set, so it can renew only the row
matching this exact client name, worker identity, session, generation and
pid.
* A terminal refusal (fenced or missing) stops it permanently and records
why. A fenced session must never beat its way back into ownership.
* A transient failure (a locked database, say) is counted and the loop
continues, so one contended write does not silently end the heartbeat.
* Nothing here raises into a caller. ``beat_once`` returns its outcome and
the thread body swallows everything, because a heartbeat failure must
degrade to "not renewed" and never crash a tool call.
"""
def __init__(
self,
registry: WorkerRegistry,
*,
worker_identity: str,
fencing_epoch: int,
session_id: str | None = None,
generation_id: str | None = None,
client_name: str | None = None,
pid: int | None = None,
ttl_seconds: float = DEFAULT_HEARTBEAT_TTL_SECONDS,
interval_seconds: float | None = None,
clock=None,
) -> None:
self._registry = registry
self.worker_identity = worker_identity
self.fencing_epoch = int(fencing_epoch)
self.session_id = session_id
self.generation_id = generation_id
self.client_name = (
normalize_client_name(client_name) if client_name is not None else None
)
self.pid = int(pid) if pid is not None else None
self.ttl_seconds = float(ttl_seconds)
self.interval_seconds = (
heartbeat_interval_for(self.ttl_seconds)
if interval_seconds is None
else heartbeat_interval_for(self.ttl_seconds, interval_seconds)
)
#: Injectable so every TTL test uses controlled time and no test waits
#: for a real interval or a real TTL to elapse.
self._clock = clock or _utc_now
self._stop_event = threading.Event()
self._thread: threading.Thread | None = None
self._lock = threading.Lock()
self._atexit_registered = False
self.started = False
self.stopped_reason: str | None = None
self.beats_attempted = 0
self.beats_renewed = 0
self.transient_failures = 0
self.last_beat_at: str | None = None
self.last_result: dict[str, Any] | None = None
# -- one beat --
def beat_once(self, now: datetime | None = None) -> dict[str, Any]:
"""Renew once. Never raises; returns the registry outcome or a failure."""
if self.stopped_reason is not None:
return {
"success": False,
"renewed": False,
"beat_attempted": False,
"reasons": [
f"supervisor already stopped: {self.stopped_reason}"
],
}
stamp = now or self._clock()
self.beats_attempted += 1
try:
result = self._registry.heartbeat(
worker_identity=self.worker_identity,
fencing_epoch=self.fencing_epoch,
now=stamp,
expected_session_id=self.session_id,
expected_generation_id=self.generation_id,
expected_client_name=self.client_name,
expected_pid=self.pid,
)
except Exception as exc:
# Transient by assumption: an unexpected error is not proof this
# worker lost ownership, so it must not silently end the heartbeat.
self.transient_failures += 1
result = {
"success": False,
"renewed": False,
"transient": True,
"blocker_kind": None,
"reasons": [f"heartbeat raised {type(exc).__name__}: {exc}"],
}
self.last_result = result
return result
self.last_result = result
if result.get("renewed"):
self.beats_renewed += 1
self.last_beat_at = result.get("last_heartbeat_at") or _ts(stamp)
return result
blocker = result.get("blocker_kind")
if blocker in TERMINAL_HEARTBEAT_BLOCKERS:
self._stop_internal(
reason=(
f"refused with blocker_kind={blocker!r}: "
+ "; ".join(result.get("reasons") or [])
)
)
else:
self.transient_failures += 1
return result
# -- lifecycle --
def start(self) -> dict[str, Any]:
"""Start the beat thread. Idempotent; safe to call from a tool call."""
with self._lock:
if self.stopped_reason is not None:
return {
"started": False,
"reasons": [f"supervisor stopped: {self.stopped_reason}"],
}
if self._thread is not None and self._thread.is_alive():
return {"started": True, "already_running": True, "reasons": []}
self._stop_event.clear()
thread = threading.Thread(
target=self._run,
name=f"gitea-worker-heartbeat-{self.worker_identity}",
daemon=True,
)
self._thread = thread
self.started = True
if not self._atexit_registered:
# Orderly shutdown stops the heartbeat. A hard kill does not
# run this, which is correct: the row must then go stale.
atexit.register(self._atexit_stop)
self._atexit_registered = True
thread.start()
return {"started": True, "already_running": False, "reasons": []}
def _run(self) -> None:
while not self._stop_event.is_set():
# Wait first: registration already stamped a fresh heartbeat, so an
# immediate beat would be a redundant write on every server launch.
if self._stop_event.wait(self.interval_seconds):
return
try:
self.beat_once()
except Exception:
# beat_once is already total; this is the last-resort guard that
# keeps a supervisor thread from dying silently.
self.transient_failures += 1
if self.stopped_reason is not None:
return
def stop(self, reason: str = "stopped") -> dict[str, Any]:
"""Stop beating. Idempotent, and prompt because the loop waits on an Event."""
self._stop_internal(reason=reason)
thread = self._thread
if thread is not None and thread.is_alive():
thread.join(timeout=max(1.0, min(5.0, self.interval_seconds)))
return {
"stopped": True,
"reason": self.stopped_reason,
"thread_alive": bool(thread is not None and thread.is_alive()),
}
def _stop_internal(self, *, reason: str) -> None:
if self.stopped_reason is None:
self.stopped_reason = reason
self._stop_event.set()
if self._atexit_registered:
try:
atexit.unregister(self._atexit_stop)
except Exception:
pass
self._atexit_registered = False
def _atexit_stop(self) -> None:
try:
self.stop(reason="process exit")
except Exception:
pass
# -- observability --
def status(self) -> dict[str, Any]:
"""Read-only observability payload; safe to embed in a tool result."""
thread = self._thread
return {
"supervised": True,
"worker_identity": self.worker_identity,
"fencing_epoch": self.fencing_epoch,
"session_id": self.session_id,
"generation_id": self.generation_id,
"client_name": self.client_name,
"pid": self.pid,
"heartbeat_ttl_seconds": self.ttl_seconds,
"heartbeat_interval_seconds": self.interval_seconds,
"started": self.started,
"running": bool(
thread is not None
and thread.is_alive()
and self.stopped_reason is None
),
"stopped_reason": self.stopped_reason,
"beats_attempted": self.beats_attempted,
"beats_renewed": self.beats_renewed,
"transient_failures": self.transient_failures,
"last_heartbeat_at": self.last_beat_at,
"last_blocker_kind": (self.last_result or {}).get("blocker_kind"),
}
+12
View File
@@ -163,6 +163,18 @@ TASK_CAPABILITY_MAP: dict[str, dict[str, str]] = {
"permission": "gitea.read",
"role": "author",
},
# #978: instance-level fleet identity/health snapshot. gitea.read is the
# operation gate; the tool additionally restricts role_kind to
# controller|reconciler so author/reviewer/merger cannot use it as a
# mutation surface and no unrelated write permission is introduced.
"snapshot_instance_fleet": {
"permission": "gitea.read",
"role": "controller",
},
"gitea_snapshot_instance_fleet": {
"permission": "gitea.read",
"role": "controller",
},
# #644: Phase 2 Web Console recovery tasks.
"clear_stale_binding": {
"permission": "gitea.read",
+14 -1
View File
@@ -175,8 +175,21 @@ class TestLauncherSnippets(unittest.TestCase):
def test_only_safe_keys_no_secrets(self):
entry = gitea_config.launcher_entry("prgs", "/cfg/profiles.json")["gitea-tools"]
self.assertEqual(set(entry), {"command", "args", "env"})
self.assertEqual(set(entry["env"]), {"GITEA_MCP_CONFIG", "GITEA_MCP_PROFILE", "GITEA_CLIENT_MANAGED"})
# #978 B1: production launcher also injects trusted client identity.
required = {
"GITEA_MCP_CONFIG",
"GITEA_MCP_PROFILE",
"GITEA_CLIENT_MANAGED",
"GITEA_MCP_CLIENT",
"GITEA_MCP_CLIENT_INSTANCE",
"GITEA_MCP_INSTANCE_PROVENANCE",
}
self.assertTrue(required.issubset(set(entry["env"])), entry["env"])
self.assertEqual(entry["env"]["GITEA_MCP_PROFILE"], "prgs")
self.assertTrue(
entry["env"]["GITEA_MCP_CLIENT_INSTANCE"].startswith("inst-"),
entry["env"]["GITEA_MCP_CLIENT_INSTANCE"],
)
blob = json.dumps(entry).lower()
for word in ("token", "password", "secret"):
self.assertNotIn(word, blob)
File diff suppressed because it is too large Load Diff
File diff suppressed because it is too large Load Diff