Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 3 additions & 3 deletions .github/actions/plugin-validate/action.yml
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
name: Hermes Plugin Validate
description: >-
Validate a Hermes Agent plugin (plugin.yaml manifest schema AND
declared-vs-actually-registered capabilities) using
declared capabilities compared with literal registration calls) using
`hermes plugins validate`. Drop this into your plugin repo's CI:

- uses: actions/checkout@<sha>
Expand Down Expand Up @@ -47,8 +47,8 @@ runs:
run: |
set -uo pipefail
# `hermes plugins validate` checks the plugin.yaml manifest schema
# and loads the plugin in a scratch subprocess to verify that the
# capabilities it DECLARES match what it actually registers.
# and statically inspects literal registration calls to compare the
# declared capabilities with those calls. Candidate code is never executed.
if hermes plugins validate "$_PLUGIN_PATH"; then
echo "✅ PASS: plugin at '$_PLUGIN_PATH' validated cleanly"
else
Expand Down
33 changes: 24 additions & 9 deletions agent/conversation_worktree.py
Original file line number Diff line number Diff line change
Expand Up @@ -498,6 +498,7 @@ def bind_new_root_session(
branch=branch,
base_commit=base_commit,
repo_common_dir=str(source_common_dir),
source_worktree=str(source),
)
except ConversationWorktreeConflict as exc:
record = self._db.get_conversation_worktree(root_session_id)
Expand Down Expand Up @@ -546,6 +547,20 @@ def bind_new_root_session(
)
raise

def _durable_source(self, record: ConversationWorktreeRecord) -> Path:
# Legacy rows can only recover from an explicitly configured checkout.
source = (
Path(record.source_worktree).resolve()
if record.source_worktree
else self._source_repository_identity()[0]
)
actual = Path(self._git_stdout(
source, ["rev-parse", "--path-format=absolute", "--git-common-dir"], "identity"
)).resolve()
if actual != Path(record.repo_common_dir).resolve():
raise ConversationWorktreeError("durable source repository identity changed", phase="identity")
return source

def resolve_existing_session(
self, root_session_id: str
) -> ConversationWorktreeBinding | None:
Expand All @@ -554,7 +569,7 @@ def resolve_existing_session(
if record is None:
return None
source_common_dir = Path(record.repo_common_dir).resolve()
source = source_common_dir.parent
source = self._durable_source(record)
path = Path(record.worktree_path).resolve()
branch = record.branch
if not source.is_dir() or not source_common_dir.is_dir():
Expand Down Expand Up @@ -647,7 +662,7 @@ def remove_after_explicit_request(

try:
source_common_dir = Path(record.repo_common_dir).resolve()
source = source_common_dir.parent
source = self._durable_source(record)
if not source.is_dir() or not source_common_dir.is_dir():
raise ConversationWorktreeError(
"durable conversation repository identity is unavailable", phase="identity"
Expand Down Expand Up @@ -756,13 +771,11 @@ def block(reason: str) -> None:

try:
source_common_dir = Path(record.repo_common_dir).resolve()
source = source_common_dir.parent
expected_path, expected_branch = Path(record.worktree_path).resolve(), record.branch
source = self._durable_source(record)
expected_path = Path(record.worktree_path).resolve()
if (
record.state not in {"ready", "retained"}
or Path(record.worktree_path).resolve() != expected_path.resolve()
or record.branch != expected_branch
or Path(record.repo_common_dir).resolve() != source_common_dir.resolve()
or not self._exact_owner_claims_present(record)
or not expected_path.is_dir()
):
return CleanupVerdict(False, ("mismatched identity",))
Expand Down Expand Up @@ -1402,8 +1415,10 @@ def _validated_ready_binding(
"ready conversation worktree no longer descends from its base commit",
phase="recovery",
)
self._ensure_common_owner_claim(record)
self._ensure_owner_marker(record)
if not self._exact_owner_claims_present(record):
raise ConversationWorktreeError(
"ready conversation worktree ownership is unavailable", phase="recovery"
)
self._event(
"conversation_worktree.reuse", root_session_id=record.root_session_id
)
Expand Down
4 changes: 2 additions & 2 deletions agent/llm_egress_runtime.py
Original file line number Diff line number Diff line change
Expand Up @@ -207,15 +207,15 @@ def egress_enforcement_enabled() -> bool:

config = load_config_readonly()
if managed_scope.is_key_managed("runtime.llm_egress_enforcement"):
posture = str((config.get("runtime") or {}).get("llm_egress_enforcement", "enabled") or "enabled")
posture = str((config.get("runtime") or {}).get("llm_egress_enforcement", "enabled"))
return posture.strip().lower() not in {"0", "false", "off", "disabled", "disable", "monitor"}
except Exception:
pass
try:
from hermes_cli.config import load_config_readonly

runtime = load_config_readonly().get("runtime") or {}
posture = str(runtime.get("llm_egress_enforcement", "enabled") or "enabled")
posture = str(runtime.get("llm_egress_enforcement", "enabled"))
return posture.strip().lower() not in {"0", "false", "off", "disabled", "disable", "monitor"}
except Exception:
# A malformed or unavailable config must not silently weaken the boundary.
Expand Down
90 changes: 52 additions & 38 deletions agent/trajectory.py
Original file line number Diff line number Diff line change
Expand Up @@ -5,8 +5,6 @@
import io
import logging
import os
import shutil
import stat
import tempfile
from datetime import datetime
from typing import Any, Dict, List
Expand Down Expand Up @@ -47,47 +45,63 @@ def _build_gzip_member(line: str) -> bytes:


def _append_gzip_member_atomically(filename: str, payload: bytes) -> None:
"""Append a complete member with a stable lock and atomic destination replace.
"""Append with a durable rollback offset, writing only the new member.

A process can be killed at any point during a regular-file write, including
between short writes. Building the new file beside the destination keeps a
killed writer from ever publishing a partial gzip member. The sidecar lock
remains stable across ``os.replace`` so concurrent writers cannot split the
critical section when the destination inode changes.
The stable lock serializes writers. After a killed writer the next append
rolls back the incomplete member before proceeding. Ordinary gzip readers
must wait for that recovery if a writer died during its destination write.
"""
directory = os.path.dirname(os.path.abspath(filename)) or "."
lock_name = f"{filename}.lock"
with open(lock_name, "a+b") as lock_file:
locked = False
journal = f"{filename}.pending"
with open(f"{filename}.lock", "a+b") as lock_file:
_lock_append_handle(lock_file, True)
try:
_lock_append_handle(lock_file, True)
locked = True
existing_mode = None
if os.path.exists(filename):
existing_mode = stat.S_IMODE(os.stat(filename).st_mode)
fd, staged_name = tempfile.mkstemp(
prefix=f".{os.path.basename(filename)}.", suffix=".tmp", dir=directory
)
try:
with os.fdopen(fd, "wb") as staged:
if os.path.exists(filename):
with open(filename, "rb") as existing:
shutil.copyfileobj(existing, staged)
staged.write(payload)
staged.flush()
os.fsync(staged.fileno())
if existing_mode is not None:
os.chmod(staged_name, existing_mode)
os.replace(staged_name, filename)
finally:
if os.path.exists(staged_name):
os.unlink(staged_name)
finally:
if locked:
with open(filename, "a+b") as destination:
if os.path.exists(journal):
with open(journal, encoding="ascii") as pending:
offset = int(pending.read())
if offset < 0 or offset > os.fstat(destination.fileno()).st_size:
raise ValueError("invalid trajectory recovery offset")
destination.truncate(offset)
destination.flush()
os.fsync(destination.fileno())
os.unlink(journal)
destination.seek(0, os.SEEK_END)
offset = destination.tell()
fd, staged_name = tempfile.mkstemp(prefix=".trajectory-", dir=directory)
try:
_lock_append_handle(lock_file, False)
except (OSError, ValueError):
pass
os.chmod(staged_name, os.fstat(destination.fileno()).st_mode & 0o777)
with os.fdopen(fd, "w", encoding="ascii") as pending:
pending.write(str(offset))
pending.flush()
os.fsync(pending.fileno())
os.replace(staged_name, journal)
_sync_trajectory_directory(directory)
try:
destination.write(payload)
destination.flush()
os.fsync(destination.fileno())
Comment thread
mrkillbob marked this conversation as resolved.
except BaseException:
destination.truncate(offset)
destination.flush()
os.fsync(destination.fileno())
raise
os.unlink(journal)
_sync_trajectory_directory(directory)
finally:
if os.path.exists(staged_name):
os.unlink(staged_name)
finally:
_lock_append_handle(lock_file, False)


def _sync_trajectory_directory(directory: str) -> None:
if os.name != "nt":
fd = os.open(directory, os.O_RDONLY)
try:
os.fsync(fd)
finally:
os.close(fd)


def save_trajectory(trajectory: List[Dict[str, Any]], model: str, completed: bool, filename: str = None):
Expand Down
5 changes: 4 additions & 1 deletion agent/vault_backends/bitwarden.py
Original file line number Diff line number Diff line change
Expand Up @@ -46,7 +46,10 @@ def _bw(self) -> Path:
return Path(found)

def _env(self, session_token: Optional[str]) -> Dict[str, str]:
env = {k: os.environ[k] for k in _ENV_KEEP if k in os.environ}
from agent.secret_scope import get_secret
env = {k: os.environ[k] for k in _ENV_KEEP if k in os.environ and k != "BITWARDENCLI_APPDATA_DIR"}
if appdata := get_secret("BITWARDENCLI_APPDATA_DIR", ""):
env["BITWARDENCLI_APPDATA_DIR"] = appdata
env["NO_COLOR"] = "1"
if session_token:
env["BW_SESSION"] = session_token
Expand Down
6 changes: 3 additions & 3 deletions agent/vault_backends/onepassword.py
Original file line number Diff line number Diff line change
Expand Up @@ -49,14 +49,14 @@ def _op(self) -> Path:

def _env(self, session_token: Optional[str]) -> Dict[str, str]:
from agent.secret_scope import get_secret
env = {k: os.environ[k] for k in _OP_ENV_ALLOWLIST if k in os.environ and not k.startswith("OP_CONNECT_")}
env = {k: os.environ[k] for k in _OP_ENV_ALLOWLIST if k in os.environ and not k.startswith("OP_")}
# Connect credentials outrank OP_SERVICE_ACCOUNT_TOKEN inside op, so they must come from the
# profile's own secret scope like the service token does — never from the launch environment.
for k in ("OP_CONNECT_HOST", "OP_CONNECT_TOKEN"):
for k in ("OP_CONNECT_HOST", "OP_CONNECT_TOKEN", "OP_LOAD_DESKTOP_APP_SETTINGS"):
if v := get_secret(k, ""):
env[k] = v
env["NO_COLOR"] = "1"
account = str(self.cfg.get("account") or "")
account = str(self.cfg.get("account") or get_secret("OP_ACCOUNT", "") or "")
if account:
env["OP_ACCOUNT"] = account
if self._service_token:
Expand Down
3 changes: 3 additions & 0 deletions apps/desktop/electron/backend-dial-claim.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -120,6 +120,7 @@ describe('backend dial routing (#90812)', () => {
it('uses the composite registry scope for a registry backend', async () => {
const claims = new BackendDialClaims()
const dial = vi.fn(async () => 'registry')

const scopeKey = vi.fn((connectionId: string | null, profile: string | null | undefined) =>
`conn:${connectionId ?? 'local'}::${profile ?? 'default'}`)

Expand All @@ -131,6 +132,7 @@ describe('backend dial routing (#90812)', () => {
it('uses the local profile scope for a local backend', async () => {
const claims = new BackendDialClaims()
const dial = vi.fn(async () => 'local')

const scopeKey = vi.fn((connectionId: string | null, profile: string | null | undefined) =>
`conn:${connectionId ?? 'local'}::${profile ?? 'default'}`)

Expand All @@ -142,6 +144,7 @@ describe('backend dial routing (#90812)', () => {
it('preserves an explicit pooled key when the parsed route would normalize it', async () => {
const claims = new BackendDialClaims()
const dial = vi.fn(async () => 'forced-local')

const scopeKey = vi.fn((connectionId: string | null, profile: string | null | undefined) =>
`conn:${connectionId ?? 'local'}::${profile ?? 'default'}`)

Expand Down
4 changes: 4 additions & 0 deletions apps/desktop/electron/desktop-background-shutdown.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -20,11 +20,14 @@ describe('Desktop background-service shutdown', () => {

it('boots out the exact Hermes companion launchd job on macOS', async () => {
const children: EventEmitter[] = []

const spawnFn = vi.fn(() => {
const child = Object.assign(new EventEmitter(), {
kill: vi.fn(() => true)
})

children.push(child)

return child
})

Expand All @@ -36,6 +39,7 @@ describe('Desktop background-service shutdown', () => {
uid: 501,
timeoutMs: 1_000
})

expect(spawnFn).toHaveBeenCalledTimes(1)
children[0].emit('exit', 0, null)

Expand Down
2 changes: 2 additions & 0 deletions apps/desktop/electron/desktop-background-shutdown.ts
Original file line number Diff line number Diff line change
Expand Up @@ -47,11 +47,13 @@ function runStopCommand(
if (settled) {
return
}

settled = true

if (timer) {
clearTimeout(timer)
}

resolve(ok)
}

Expand Down
5 changes: 5 additions & 0 deletions apps/desktop/electron/dispatcher-readiness.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,7 @@ test('runDispatcherReadinessGate advances the boot phase before checking readine
'session-token',
async () => {
events.push('readiness-checked')

return { status: 'ready', ready: true, gateway_pid: 1, message: 'ok' }
},
async (id: string) => {
Expand Down Expand Up @@ -41,6 +42,7 @@ test('accepts a live gateway-owned dispatcher', async () => {

test('starts one supervised gateway when the dispatcher is offline and waits for readiness', async () => {
const calls: Array<[string, string | null, string]> = []

const responses = [
{ status: 'offline', ready: false, gateway_pid: null, message: 'gateway is offline' },
{ ok: true, pid: 7654, name: 'gateway-start' },
Expand All @@ -53,6 +55,7 @@ test('starts one supervised gateway when the dispatcher is offline and waits for
'session-token',
async (url, token, options = {}) => {
calls.push([url, token, options.method || 'GET'])

return responses.shift()
},
{ attempts: 2, pollMs: 0, sleep: async () => {} }
Expand All @@ -72,6 +75,7 @@ test('allows Desktop startup when the embedded dispatcher is disabled', async ()

const result = await ensureKanbanDispatcherReady('http://127.0.0.1:9000', 'session-token', async url => {
calls.push(url)

return { status: 'disabled', ready: false, gateway_pid: null, message: 'dispatcher is disabled' }
})

Expand All @@ -98,6 +102,7 @@ test('blocks startup without starting a gateway for unknown dispatcher state', a
await assert.rejects(
ensureKanbanDispatcherReady('http://127.0.0.1:9000', 'session-token', async url => {
calls.push(url)

return { status: 'unknown', ready: false, gateway_pid: null, message: 'dispatcher is unknown' }
}),
error => {
Expand Down
4 changes: 4 additions & 0 deletions apps/desktop/electron/dispatcher-readiness.ts
Original file line number Diff line number Diff line change
Expand Up @@ -62,9 +62,11 @@ export async function ensureKanbanDispatcherReady(
// Treat it the same as an explicit { status: "disabled" } response —
// Desktop startup must not fail because an optional plugin is turned off.
const detail = error instanceof Error ? error.message : String(error)

if (/^404[^\d]/.test(detail)) {
return { status: 'disabled', ready: false, gateway_pid: null, message: detail }
}

throw new DispatcherReadinessError(`dispatcher readiness could not be verified: ${detail}`)
}

Expand Down Expand Up @@ -106,6 +108,7 @@ export async function ensureKanbanDispatcherReady(
const detail = error instanceof Error ? error.message : String(error)
throw new DispatcherReadinessError(`dispatcher readiness could not be verified after gateway start: ${detail}`)
}

continue
}

Expand Down Expand Up @@ -133,5 +136,6 @@ export async function runDispatcherReadinessGate(
advancePhase: AdvancePhase
): Promise<DispatcherReadiness> {
await advancePhase('backend.dispatcher', 'Verifying Kanban dispatcher readiness', 92)

return ensureKanbanDispatcherReady(baseUrl, token, fetchJson)
}
1 change: 1 addition & 0 deletions apps/desktop/electron/main-window-lifecycle.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@ test('closing the last desktop window quits on every platform', () => {
test('destroying renderer windows for a deferred drain is idempotent', () => {
let liveDestroyed = 0
let alreadyDestroyed = 0

const windows = [
{ destroy: () => void (liveDestroyed += 1), isDestroyed: () => false },
{ destroy: () => void (alreadyDestroyed += 1), isDestroyed: () => true }
Expand Down
Loading
Loading