From 01c58ad0fe8a9a99f5a3878714ec4584e1a1c60a Mon Sep 17 00:00:00 2001 From: Harbor404 <2657212322@qq.com> Date: Thu, 1 Oct 2026 16:49:47 +0800 Subject: [PATCH 1/3] feat(source-cache): add configuration and presigner Part of #2732 (CP1.1). - add off-by-default sourceCache Helm values with member-only delivery - validate daemon configuration and add SigV4 GET/PUT presigning - cover chart rendering, config validation, signature contracts, and MinIO --- .../agentconnect/templates/daemon-pool.yaml | 55 +++++ charts/agentconnect/values.yaml | 31 +++ packages/daemon/src/config/config-schema.ts | 4 + packages/daemon/src/config/load-config.ts | 8 + packages/daemon/src/source-cache/config.ts | 163 ++++++++++++++ .../daemon/src/source-cache/credentials.ts | 198 ++++++++++++++++ packages/daemon/src/source-cache/signer.ts | 212 ++++++++++++++++++ .../daemon/test/source-cache-config.test.ts | 199 ++++++++++++++++ .../daemon/test/source-cache-minio.it.test.ts | 63 ++++++ .../daemon/test/source-cache-signer.test.ts | 126 +++++++++++ scripts/test-chart-render.rb | 90 ++++++++ 11 files changed, 1149 insertions(+) create mode 100644 packages/daemon/src/source-cache/config.ts create mode 100644 packages/daemon/src/source-cache/credentials.ts create mode 100644 packages/daemon/src/source-cache/signer.ts create mode 100644 packages/daemon/test/source-cache-config.test.ts create mode 100644 packages/daemon/test/source-cache-minio.it.test.ts create mode 100644 packages/daemon/test/source-cache-signer.test.ts diff --git a/charts/agentconnect/templates/daemon-pool.yaml b/charts/agentconnect/templates/daemon-pool.yaml index 88c6daa46..3e92bd558 100644 --- a/charts/agentconnect/templates/daemon-pool.yaml +++ b/charts/agentconnect/templates/daemon-pool.yaml @@ -2,6 +2,36 @@ {{- $dataPlaneSecret := required "daemonPool.dataPlane.existingSecret is required when daemonPool.enabled" .Values.daemonPool.dataPlane.existingSecret }} {{- $sandboxNamespace := .Values.daemonPool.sandboxNamespace | default .Release.Namespace }} {{- $daemonTag := .Values.daemonPool.tag | default (include "agentconnect.imageTag" .) }} +{{- $sourceCache := .Values.sourceCache }} +{{- $sourceCacheEndpoint := trim ($sourceCache.endpoint | default "") }} +{{- $sourceCacheBucket := trim ($sourceCache.bucket | default "") }} +{{- $sourceCacheEnabled := and (ne $sourceCacheEndpoint "") (ne $sourceCacheBucket "") }} +{{- if and (ne $sourceCacheEndpoint "") (eq $sourceCacheBucket "") }} +{{- fail "sourceCache.bucket is required when sourceCache.endpoint is set" }} +{{- end }} +{{- if and (eq $sourceCacheEndpoint "") (ne $sourceCacheBucket "") }} +{{- fail "sourceCache.endpoint is required when sourceCache.bucket is set" }} +{{- end }} +{{- if $sourceCacheEnabled }} +{{- if not (has $sourceCache.credentialSource (list "serviceAccount" "secret")) }} +{{- fail "sourceCache.credentialSource must be serviceAccount or secret" }} +{{- end }} +{{- if eq $sourceCache.credentialSource "secret" }} +{{- if not $sourceCache.existingSecret }} +{{- fail "sourceCache.existingSecret is required when sourceCache.credentialSource=secret" }} +{{- end }} +{{- if or (not $sourceCache.accessKeyIdKey) (not $sourceCache.secretAccessKeyKey) }} +{{- fail "sourceCache.accessKeyIdKey and sourceCache.secretAccessKeyKey are required for Secret credentials" }} +{{- end }} +{{- else if $sourceCache.existingSecret }} +{{- fail "sourceCache.existingSecret is only valid when sourceCache.credentialSource=secret" }} +{{- end }} +{{- range $name := list "AC_SOURCE_CACHE_CONFIG" "AC_SOURCE_CACHE_ACCESS_KEY_ID" "AC_SOURCE_CACHE_SECRET_ACCESS_KEY" "AC_SOURCE_CACHE_SESSION_TOKEN" }} +{{- if hasKey $.Values.daemonPool.extraEnv $name }} +{{- fail (printf "daemonPool.extraEnv.%s collides with the chart-owned Source Cache configuration" $name) }} +{{- end }} +{{- end }} +{{- end }} # The install-wide daemon pool. Its Pods are pool members. apiVersion: v1 kind: ServiceAccount @@ -165,6 +195,31 @@ spec: {{- include "agentconnect.daemonPoolClusterEnv" . | nindent 12 }} - name: AC_K8S_SHIM_PORT value: {{ .Values.daemonPool.shimPort | quote }} + {{- if $sourceCacheEnabled }} + # Member-only Source Cache configuration. Credential values are always projected + # by Secret reference; the static key path below never renders a value inline. + - name: AC_SOURCE_CACHE_CONFIG + value: {{ dict "endpoint" $sourceCacheEndpoint "region" $sourceCache.region "bucket" $sourceCacheBucket "prefix" $sourceCache.prefix "forcePathStyle" $sourceCache.forcePathStyle "credentialSource" $sourceCache.credentialSource "limits" $sourceCache.limits | toJson | quote }} + {{- if eq $sourceCache.credentialSource "secret" }} + - name: AC_SOURCE_CACHE_ACCESS_KEY_ID + valueFrom: + secretKeyRef: + name: {{ $sourceCache.existingSecret | quote }} + key: {{ $sourceCache.accessKeyIdKey | quote }} + - name: AC_SOURCE_CACHE_SECRET_ACCESS_KEY + valueFrom: + secretKeyRef: + name: {{ $sourceCache.existingSecret | quote }} + key: {{ $sourceCache.secretAccessKeyKey | quote }} + {{- if $sourceCache.sessionTokenKey }} + - name: AC_SOURCE_CACHE_SESSION_TOKEN + valueFrom: + secretKeyRef: + name: {{ $sourceCache.existingSecret | quote }} + key: {{ $sourceCache.sessionTokenKey | quote }} + {{- end }} + {{- end }} + {{- end }} # The sink the readiness probe below reads: the daemon serves GET /readyz here and # answers 200 only once startup finished, the control plane acknowledged registration # and the install-wide runtime probe returned — 503 otherwise, including from SIGTERM diff --git a/charts/agentconnect/values.yaml b/charts/agentconnect/values.yaml index 1796f241a..13be27cca 100644 --- a/charts/agentconnect/values.yaml +++ b/charts/agentconnect/values.yaml @@ -409,6 +409,37 @@ daemonPool: nodeSelector: {} tolerations: [] +# Optional Source Cache. Empty endpoint or bucket leaves it OFF: no URLs are issued and +# pool members clone from origin exactly as before. Only the daemon-pool member receives +# this configuration; sandbox pods never receive bucket credentials. +sourceCache: + # Full S3 API origin, e.g. https://s3.us-east-1.amazonaws.com or https://minio.example.test + endpoint: '' + # SigV4 region. Use "auto" for R2; a concrete region for S3 or MinIO. + region: auto + # Bucket (or bucket prefix) dedicated to this install's Source Cache. + bucket: '' + # Optional key prefix inside the bucket. The daemon adds src/, files/, and snapshots/ below it. + prefix: '' + # MinIO and other S3-compatible stores commonly require path-style addressing. + forcePathStyle: false + # serviceAccount uses the pool member's IAM/Workload Identity. secret projects static keys + # from an existing Secret by name; it never embeds credential values in this values file. + credentialSource: serviceAccount + existingSecret: '' + accessKeyIdKey: AWS_ACCESS_KEY_ID + secretAccessKeyKey: AWS_SECRET_ACCESS_KEY + # Optional; set to "" for stores or identities that do not issue a session token. + sessionTokenKey: AWS_SESSION_TOKEN + # Section 10 limits. The daemon validates the values; the data plane enforces them. + limits: + maxBundleBytes: 2147483648 # 2 GiB + orgTotalBytes: 21474836480 # 20 GiB + pendingReservationMs: 3600000 # 1 h + pendingObjectTtlDays: 2 + unreferencedObjectTtlDays: 7 + unreadPointerTtlDays: 30 + controlPlane: image: control-plane # Per-component tag override; empty ⇒ the shared image.tag. Set it to the diff --git a/packages/daemon/src/config/config-schema.ts b/packages/daemon/src/config/config-schema.ts index 565fff55f..2f1c48162 100644 --- a/packages/daemon/src/config/config-schema.ts +++ b/packages/daemon/src/config/config-schema.ts @@ -6,6 +6,7 @@ import { SESSION_RETENTION_RE, type SessionRetentionSetting } from '@agentconnect.md/protocol' +import { SourceCacheConfigSchema } from '../source-cache/config.js' /** The `{name, value}[]` shape shared by runtime env, MCP env, and MCP headers. */ const NameValueList = z.array(z.object({ name: z.string(), value: z.string() })).default([]) @@ -246,6 +247,9 @@ export const ConfigSchema = z.object({ // degradation); the CP's register/ok snapshot re-converges it authoritatively once // connected. CP-owned — overwritten on every roster converge, not hand-edited. relays: z.array(RelayRosterEntry).default([]), + // Optional Source Cache. Absence is the disabled no-op; the environment document is + // parsed by loadConfig and never exposed to sandboxes or the control plane. + sourceCache: SourceCacheConfigSchema.optional(), logging: z .object({ level: z.enum(['trace', 'debug', 'info', 'warn', 'error']).default('info') }) .default({ level: 'info' }), diff --git a/packages/daemon/src/config/load-config.ts b/packages/daemon/src/config/load-config.ts index 6aac7de72..9b600814e 100644 --- a/packages/daemon/src/config/load-config.ts +++ b/packages/daemon/src/config/load-config.ts @@ -8,6 +8,7 @@ import { } from '@agentconnect.md/protocol' import { z } from 'zod' import { ConfigSchema, WorkspaceGitOrigin, type Config } from './config-schema.js' +import { sourceCacheConfigFromEnv } from '../source-cache/config.js' import { resolveRoot, configPath, defaultAgentsDir } from '../paths.js' export interface FlatOverrides { @@ -93,6 +94,13 @@ export function loadConfig( throw new Error(`config not found: ${file} (create it, pass --config, or run \`agentconnect login\`)`) } const cfg = ConfigSchema.parse(raw) // throws on invalid + // Pool members have no writable config file. The chart supplies the Source Cache + // document by environment; a file that states sourceCache wins (same precedence as + // the other operator-owned policy overrides below). + if ((raw as { sourceCache?: unknown } | null)?.sourceCache === undefined) { + const sourceCache = sourceCacheConfigFromEnv(process.env) + if (sourceCache) cfg.sourceCache = sourceCache + } for (const warning of mapLegacySandbox(cfg)) opts.warn?.(warning) const o = opts.overrides ?? {} diff --git a/packages/daemon/src/source-cache/config.ts b/packages/daemon/src/source-cache/config.ts new file mode 100644 index 000000000..9924011da --- /dev/null +++ b/packages/daemon/src/source-cache/config.ts @@ -0,0 +1,163 @@ +import { z } from 'zod' + +/** The one environment document the Helm chart renders into a pool member. */ +export const SOURCE_CACHE_CONFIG_ENV = 'AC_SOURCE_CACHE_CONFIG' +/** Static S3 credentials are projected by Secret reference, never inline in the config document. */ +export const SOURCE_CACHE_ACCESS_KEY_ID_ENV = 'AC_SOURCE_CACHE_ACCESS_KEY_ID' +export const SOURCE_CACHE_SECRET_ACCESS_KEY_ENV = 'AC_SOURCE_CACHE_SECRET_ACCESS_KEY' +export const SOURCE_CACHE_SESSION_TOKEN_ENV = 'AC_SOURCE_CACHE_SESSION_TOKEN' + +const GIB = 1024 ** 3 +const HOUR_MS = 60 * 60 * 1000 + +export const SOURCE_CACHE_LIMIT_DEFAULTS = { + maxBundleBytes: 2 * GIB, + orgTotalBytes: 20 * GIB, + pendingReservationMs: HOUR_MS, + pendingObjectTtlDays: 2, + unreferencedObjectTtlDays: 7, + unreadPointerTtlDays: 30 +} as const + +/** Section 10 defaults. The data plane owns enforcement; config validation rejects values + * that cannot be represented or would invert an invariant. */ +export const SourceCacheLimitsSchema = z + .object({ + maxBundleBytes: z.number().int().positive().default(SOURCE_CACHE_LIMIT_DEFAULTS.maxBundleBytes), + orgTotalBytes: z.number().int().positive().default(SOURCE_CACHE_LIMIT_DEFAULTS.orgTotalBytes), + pendingReservationMs: z.number().int().positive().default(SOURCE_CACHE_LIMIT_DEFAULTS.pendingReservationMs), + pendingObjectTtlDays: z.number().int().positive().default(SOURCE_CACHE_LIMIT_DEFAULTS.pendingObjectTtlDays), + unreferencedObjectTtlDays: z + .number() + .int() + .positive() + .default(SOURCE_CACHE_LIMIT_DEFAULTS.unreferencedObjectTtlDays), + unreadPointerTtlDays: z.number().int().positive().default(SOURCE_CACHE_LIMIT_DEFAULTS.unreadPointerTtlDays) + }) + .strict() + +const BucketName = z + .string() + .min(3) + .max(63) + .regex(/^[a-z0-9][a-z0-9.-]*[a-z0-9]$/, 'invalid S3 bucket name') + .refine((value) => !value.includes('..') && !value.includes('.-') && !value.includes('-.'), 'invalid S3 bucket name') + +const Prefix = z + .string() + .max(1024) + .transform((value, ctx) => { + const normalized = value.replace(/^\/+|\/+$/g, '') + if ( + (normalized !== '' && + normalized + .split('/') + .some((segment) => segment === '' || segment === '.' || segment === '..' || segment.includes('\\'))) || + /[\0-\x1f\x7f]/.test(normalized) + ) { + ctx.addIssue({ code: 'custom', message: 'invalid source cache prefix' }) + return z.NEVER + } + return normalized + }) + +const Endpoint = z + .string() + .url() + .transform((value, ctx) => { + const url = new URL(value) + if (url.username || url.password || url.search || url.hash || url.pathname !== '/') { + ctx.addIssue({ + code: 'custom', + message: 'source cache endpoint must be an origin without credentials, path, query, or fragment' + }) + return z.NEVER + } + return url.origin + }) + +const ServiceAccountCredentials = z.object({ source: z.literal('serviceAccount') }).strict() +const SecretCredentials = z + .object({ + source: z.literal('secret'), + accessKeyId: z + .string() + .min(1) + .max(16 * 1024), + secretAccessKey: z + .string() + .min(1) + .max(16 * 1024), + sessionToken: z + .string() + .min(1) + .max(64 * 1024) + .optional() + }) + .strict() + +export const SourceCacheCredentialsSchema = z.discriminatedUnion('source', [ + ServiceAccountCredentials, + SecretCredentials +]) + +const SourceCacheBaseSchema = z + .object({ + endpoint: Endpoint, + region: z.string().trim().min(1).max(64).default('auto'), + bucket: BucketName, + prefix: Prefix.default(''), + forcePathStyle: z.boolean().default(false), + limits: SourceCacheLimitsSchema.default(SOURCE_CACHE_LIMIT_DEFAULTS) + }) + .strict() + +/** The daemon's validated shape. Presence means the feature is enabled; absence is the no-op. */ +export const SourceCacheConfigSchema = SourceCacheBaseSchema.extend({ + credentials: SourceCacheCredentialsSchema +}).strict() + +export type SourceCacheConfig = z.infer +export type SourceCacheCredentials = SourceCacheConfig['credentials'] + +const SourceCacheEnvironmentSchema = SourceCacheBaseSchema.extend({ + credentialSource: z.enum(['serviceAccount', 'secret']) +}).strict() + +function requiredSecret(env: NodeJS.ProcessEnv, name: string): string { + const value = env[name]?.trim() + if (!value) throw new Error(`${name} is required when source cache credentialSource=secret`) + return value +} + +/** Parse the member-only Helm document. Blank means the feature is off; malformed means + * configuration is wrong and must be fixed rather than silently disabling the cache. */ +export function sourceCacheConfigFromEnv(env: NodeJS.ProcessEnv = process.env): SourceCacheConfig | undefined { + const raw = env[SOURCE_CACHE_CONFIG_ENV]?.trim() + if (!raw) return undefined + + let parsed: unknown + try { + parsed = JSON.parse(raw) + } catch { + throw new Error(`${SOURCE_CACHE_CONFIG_ENV} must be valid JSON`) + } + + const document = SourceCacheEnvironmentSchema.safeParse(parsed) + if (!document.success) throw new Error(`${SOURCE_CACHE_CONFIG_ENV}: ${document.error.message}`) + + const credentials = + document.data.credentialSource === 'secret' + ? { + source: 'secret' as const, + accessKeyId: requiredSecret(env, SOURCE_CACHE_ACCESS_KEY_ID_ENV), + secretAccessKey: requiredSecret(env, SOURCE_CACHE_SECRET_ACCESS_KEY_ENV), + ...(env[SOURCE_CACHE_SESSION_TOKEN_ENV]?.trim() + ? { sessionToken: env[SOURCE_CACHE_SESSION_TOKEN_ENV]!.trim() } + : {}) + } + : { source: 'serviceAccount' as const } + + const { credentialSource: _credentialSource, ...base } = document.data + return SourceCacheConfigSchema.parse({ ...base, credentials }) +} diff --git a/packages/daemon/src/source-cache/credentials.ts b/packages/daemon/src/source-cache/credentials.ts new file mode 100644 index 000000000..5aa928464 --- /dev/null +++ b/packages/daemon/src/source-cache/credentials.ts @@ -0,0 +1,198 @@ +import { readFile } from 'node:fs/promises' +import type { SourceCacheCredentials } from './config.js' + +export interface AwsCredentials { + accessKeyId: string + secretAccessKey: string + sessionToken?: string + expiration?: number +} + +export type AwsCredentialsProvider = () => Promise + +export interface CredentialDependencies { + env?: NodeJS.ProcessEnv + fetch?: typeof fetch + now?: () => Date +} + +const FIVE_MINUTES_MS = 5 * 60 * 1000 +const IMDS_CONTAINER_HOSTS = new Set(['127.0.0.1', 'localhost', '169.254.170.2', '169.254.170.23']) + +function isRecord(value: unknown): value is Record { + return typeof value === 'object' && value !== null && !Array.isArray(value) +} + +function textField(value: unknown, name: string): string { + if (!isRecord(value) || typeof value[name] !== 'string' || !value[name]) { + throw new Error(`source cache credential response is missing ${name}`) + } + return value[name] +} + +function expirationMs(value: unknown): number | undefined { + if (typeof value !== 'string' || !value) return undefined + const parsed = Date.parse(value) + return Number.isFinite(parsed) ? parsed : undefined +} + +function xmlText(xml: string, name: string): string { + const match = new RegExp(`<${name}>([\\s\\S]*?)`).exec(xml) + if (!match) throw new Error(`source cache STS response is missing ${name}`) + return (match[1] ?? '') + .replace(/</g, '<') + .replace(/>/g, '>') + .replace(/"/g, '"') + .replace(/'/g, "'") + .replace(/&/g, '&') +} + +function staticProvider(credentials: AwsCredentials): AwsCredentialsProvider { + return async () => credentials +} + +function cached(provider: AwsCredentialsProvider, now: () => Date): AwsCredentialsProvider { + let current: AwsCredentials | undefined + let inFlight: Promise | undefined + + return async () => { + const nowMs = now().getTime() + if (current && (current.expiration === undefined || current.expiration - nowMs > FIVE_MINUTES_MS)) return current + if (!inFlight) { + inFlight = provider() + .then((credentials) => { + current = credentials + return credentials + }) + .finally(() => { + inFlight = undefined + }) + } + return inFlight + } +} + +function environmentCredentials(env: NodeJS.ProcessEnv): AwsCredentials | undefined { + const accessKeyId = env.AWS_ACCESS_KEY_ID?.trim() + const secretAccessKey = env.AWS_SECRET_ACCESS_KEY?.trim() + if (!accessKeyId || !secretAccessKey) return undefined + const sessionToken = env.AWS_SESSION_TOKEN?.trim() + return { accessKeyId, secretAccessKey, ...(sessionToken ? { sessionToken } : {}) } +} + +async function webIdentityCredentials( + env: NodeJS.ProcessEnv, + region: string, + fetchImpl: typeof fetch +): Promise { + const roleArn = env.AWS_ROLE_ARN?.trim() + const tokenFile = env.AWS_WEB_IDENTITY_TOKEN_FILE?.trim() + if (!roleArn || !tokenFile) return undefined + + const token = (await readFile(tokenFile, 'utf8')).trim() + if (!token) throw new Error('source cache web identity token file is empty') + const sessionName = env.AWS_ROLE_SESSION_NAME?.trim() || 'agentconnect-source-cache' + const host = region === 'auto' ? 'sts.amazonaws.com' : `sts.${region}.amazonaws.com` + const endpoint = env.AWS_STS_ENDPOINT?.trim() || `https://${host}` + const body = new URLSearchParams({ + Action: 'AssumeRoleWithWebIdentity', + Version: '2011-06-15', + RoleArn: roleArn, + RoleSessionName: sessionName, + WebIdentityToken: token + }) + + const response = await fetchImpl(endpoint, { + method: 'POST', + headers: { 'content-type': 'application/x-www-form-urlencoded' }, + body + }) + const text = await response.text() + if (!response.ok) throw new Error(`source cache STS AssumeRoleWithWebIdentity failed with HTTP ${response.status}`) + + const expiration = Date.parse(xmlText(text, 'Expiration')) + return { + accessKeyId: xmlText(text, 'AccessKeyId'), + secretAccessKey: xmlText(text, 'SecretAccessKey'), + sessionToken: xmlText(text, 'SessionToken'), + ...(Number.isFinite(expiration) ? { expiration } : {}) + } +} + +async function authorizationToken(env: NodeJS.ProcessEnv): Promise { + const direct = env.AWS_CONTAINER_AUTHORIZATION_TOKEN?.trim() + if (direct) return direct + const file = env.AWS_CONTAINER_AUTHORIZATION_TOKEN_FILE?.trim() + if (!file) return undefined + const value = (await readFile(file, 'utf8')).trim() + return value || undefined +} + +async function containerCredentials( + env: NodeJS.ProcessEnv, + fetchImpl: typeof fetch +): Promise { + const relative = env.AWS_CONTAINER_CREDENTIALS_RELATIVE_URI?.trim() + const full = env.AWS_CONTAINER_CREDENTIALS_FULL_URI?.trim() + if (!relative && !full) return undefined + + let endpoint: URL + if (relative) { + endpoint = new URL(relative, 'http://169.254.170.2') + if (endpoint.hostname !== '169.254.170.2') throw new Error('invalid source cache container credential relative URI') + } else { + endpoint = new URL(full!) + if ( + !['http:', 'https:'].includes(endpoint.protocol) || + (endpoint.protocol === 'http:' && !IMDS_CONTAINER_HOSTS.has(endpoint.hostname)) + ) { + throw new Error('source cache container credential URI is not an allowed metadata endpoint') + } + } + + const token = await authorizationToken(env) + const response = await fetchImpl(endpoint, { + headers: token ? { authorization: token } : undefined + }) + if (!response.ok) throw new Error(`source cache container credential endpoint failed with HTTP ${response.status}`) + const payload: unknown = await response.json() + return { + accessKeyId: textField(payload, 'AccessKeyId'), + secretAccessKey: textField(payload, 'SecretAccessKey'), + sessionToken: textField(payload, 'Token'), + expiration: expirationMs(isRecord(payload) ? payload.Expiration : undefined) + } +} + +async function serviceAccountCredentials( + region: string, + env: NodeJS.ProcessEnv, + fetchImpl: typeof fetch +): Promise { + const fromEnvironment = environmentCredentials(env) + if (fromEnvironment) return fromEnvironment + const fromWebIdentity = await webIdentityCredentials(env, region, fetchImpl) + if (fromWebIdentity) return fromWebIdentity + const fromContainer = await containerCredentials(env, fetchImpl) + if (fromContainer) return fromContainer + throw new Error('source cache serviceAccount credentials are unavailable') +} + +export function sourceCacheCredentialsProvider( + credentials: SourceCacheCredentials, + region: string, + deps: CredentialDependencies = {} +): AwsCredentialsProvider { + if (credentials.source === 'secret') { + return staticProvider({ + accessKeyId: credentials.accessKeyId, + secretAccessKey: credentials.secretAccessKey, + ...(credentials.sessionToken ? { sessionToken: credentials.sessionToken } : {}) + }) + } + + const env = deps.env ?? process.env + const fetchImpl = deps.fetch ?? fetch + const now = deps.now ?? (() => new Date()) + return cached(() => serviceAccountCredentials(region, env, fetchImpl), now) +} diff --git a/packages/daemon/src/source-cache/signer.ts b/packages/daemon/src/source-cache/signer.ts new file mode 100644 index 000000000..412db2a5a --- /dev/null +++ b/packages/daemon/src/source-cache/signer.ts @@ -0,0 +1,212 @@ +import { createHash, createHmac } from 'node:crypto' +import { isIP } from 'node:net' +import { SourceCacheConfigSchema, type SourceCacheConfig } from './config.js' +import { + sourceCacheCredentialsProvider, + type AwsCredentials, + type AwsCredentialsProvider, + type CredentialDependencies +} from './credentials.js' + +export const SOURCE_CACHE_GET_TTL_SECONDS = 5 * 60 +export const SOURCE_CACHE_PUT_TTL_SECONDS = 15 * 60 +export const SOURCE_CACHE_MAX_PRESIGN_SECONDS = 7 * 24 * 60 * 60 + +const UNSIGNED_PAYLOAD = 'UNSIGNED-PAYLOAD' +const ALGORITHM = 'AWS4-HMAC-SHA256' +const SERVICE = 's3' +const TERMINATOR = 'aws4_request' +const PENDING_TAG = 'ac-cache=pending' + +export interface PresignOptions { + ttlSeconds?: number +} + +export interface PresignPutInput extends PresignOptions { + contentLength: number + checksumSha256: string +} + +export interface SourceCacheSigner { + presignGet(key: string, options?: PresignOptions): Promise + presignPut(key: string, input: PresignPutInput): Promise +} + +export interface SourceCacheSignerDependencies extends CredentialDependencies { + credentialsProvider?: AwsCredentialsProvider +} + +function sha256Hex(value: string): string { + return createHash('sha256').update(value, 'utf8').digest('hex') +} + +function hmac(key: Buffer | string, value: string): Buffer { + return createHmac('sha256', key).update(value, 'utf8').digest() +} + +function encodeRfc3986(value: string): string { + return encodeURIComponent(value).replace(/[!'()*]/g, (char) => `%${char.charCodeAt(0).toString(16).toUpperCase()}`) +} + +function canonicalPath(path: string): string { + return path + .split('/') + .map((segment) => encodeRfc3986(segment)) + .join('/') +} + +function canonicalQuery(parameters: ReadonlyArray): string { + return parameters + .map(([name, value]) => [encodeRfc3986(name), encodeRfc3986(value)] as const) + .sort(([leftName, leftValue], [rightName, rightValue]) => { + if (leftName !== rightName) return leftName < rightName ? -1 : 1 + if (leftValue === rightValue) return 0 + return leftValue < rightValue ? -1 : 1 + }) + .map(([name, value]) => `${name}=${value}`) + .join('&') +} + +function formatAmzDate(date: Date): { timestamp: string; day: string } { + if (!Number.isFinite(date.getTime())) throw new Error('source cache signer clock returned an invalid date') + const timestamp = date.toISOString().replace(/[:-]|\.\d{3}/g, '') + return { timestamp, day: timestamp.slice(0, 8) } +} + +function ttlSeconds(value: number | undefined, fallback: number): number { + const ttl = value ?? fallback + if (!Number.isSafeInteger(ttl) || ttl <= 0) throw new Error('source cache presign TTL must be a positive integer') + if (ttl > SOURCE_CACHE_MAX_PRESIGN_SECONDS) { + throw new Error( + `source cache presign TTL cannot exceed the SigV4 seven-day ceiling (${SOURCE_CACHE_MAX_PRESIGN_SECONDS}s)` + ) + } + return ttl +} + +function objectKey(prefix: string, key: string): string { + if (!key || key !== key.trim() || key.startsWith('/') || key.endsWith('/')) + throw new Error('invalid source cache object key') + const segments = key.split('/') + if ( + segments.some( + (segment) => + !segment || segment === '.' || segment === '..' || segment.includes('\\') || /[\0-\x1f\x7f]/.test(segment) + ) + ) { + throw new Error('invalid source cache object key') + } + return prefix ? `${prefix}/${key}` : key +} + +function checksumSha256(value: string): string { + const normalized = value.trim() + const decoded = Buffer.from(normalized, 'base64') + if (decoded.length !== 32 || decoded.toString('base64') !== normalized) { + throw new Error('source cache PUT checksum must be a base64-encoded SHA-256 digest') + } + return normalized +} + +function contentLength(value: number): string { + if (!Number.isSafeInteger(value) || value <= 0) + throw new Error('source cache PUT Content-Length must be a positive integer') + return String(value) +} + +function signingKey(secretAccessKey: string, day: string, region: string): Buffer { + return hmac(hmac(hmac(hmac(`AWS4${secretAccessKey}`, day), region), SERVICE), TERMINATOR) +} + +class AwsSigV4SourceCacheSigner implements SourceCacheSigner { + private readonly endpoint: URL + + constructor( + private readonly config: SourceCacheConfig, + private readonly credentials: AwsCredentialsProvider, + private readonly now: () => Date + ) { + this.endpoint = new URL(config.endpoint) + } + + async presignGet(key: string, options: PresignOptions = {}): Promise { + return this.presign('GET', key, {}, ttlSeconds(options.ttlSeconds, SOURCE_CACHE_GET_TTL_SECONDS)) + } + + async presignPut(key: string, input: PresignPutInput): Promise { + return this.presign( + 'PUT', + key, + { + 'content-length': contentLength(input.contentLength), + 'x-amz-checksum-sha256': checksumSha256(input.checksumSha256), + 'x-amz-tagging': PENDING_TAG + }, + ttlSeconds(input.ttlSeconds, SOURCE_CACHE_PUT_TTL_SECONDS) + ) + } + + private async presign( + method: 'GET' | 'PUT', + key: string, + headers: Record, + expires: number + ): Promise { + const credentials: AwsCredentials = await this.credentials() + const fullKey = objectKey(this.config.prefix, key) + const request = this.objectRequest(fullKey) + const { timestamp, day } = formatAmzDate(this.now()) + const scope = `${day}/${this.config.region}/${SERVICE}/${TERMINATOR}` + const canonicalHeaders = { + host: request.host, + ...headers + } + const signedHeaderNames = Object.keys(canonicalHeaders).sort() + const signedHeaders = signedHeaderNames.join(';') + const parameters: Array = [ + ['X-Amz-Algorithm', ALGORITHM], + ['X-Amz-Credential', `${credentials.accessKeyId}/${scope}`], + ['X-Amz-Date', timestamp], + ['X-Amz-Expires', String(expires)], + ['X-Amz-SignedHeaders', signedHeaders], + ...(credentials.sessionToken ? ([['X-Amz-Security-Token', credentials.sessionToken]] as const) : []) + ] + const query = canonicalQuery(parameters) + const canonicalHeaderBlock = signedHeaderNames + .map((name) => `${name}:${canonicalHeaders[name as keyof typeof canonicalHeaders].trim().replace(/\s+/g, ' ')}\n`) + .join('') + const canonicalRequest = [method, request.path, query, canonicalHeaderBlock, signedHeaders, UNSIGNED_PAYLOAD].join( + '\n' + ) + const stringToSign = [ALGORITHM, timestamp, scope, sha256Hex(canonicalRequest)].join('\n') + const signature = createHmac('sha256', signingKey(credentials.secretAccessKey, day, this.config.region)) + .update(stringToSign, 'utf8') + .digest('hex') + + return `${this.endpoint.protocol}//${request.host}${request.path}?${query}&X-Amz-Signature=${signature}` + } + + private objectRequest(key: string): { host: string; path: string } { + const ipLiteral = this.endpoint.hostname.startsWith('[') || isIP(this.endpoint.hostname) !== 0 + const pathStyle = this.config.forcePathStyle || ipLiteral || this.config.bucket.includes('.') + const path = pathStyle ? `/${this.config.bucket}/${key}` : `/${key}` + return { + host: pathStyle ? this.endpoint.host : `${this.config.bucket}.${this.endpoint.host}`, + path: canonicalPath(path) + } + } +} + +/** The disabled default allocates no signer, credentials provider, or network state. */ +export function createSourceCacheSigner( + config: SourceCacheConfig | undefined, + deps: SourceCacheSignerDependencies = {} +): SourceCacheSigner | undefined { + if (!config) return undefined + const parsed = SourceCacheConfigSchema.parse(config) + const now = deps.now ?? (() => new Date()) + const credentials = + deps.credentialsProvider ?? + sourceCacheCredentialsProvider(parsed.credentials, parsed.region, { env: deps.env, fetch: deps.fetch, now }) + return new AwsSigV4SourceCacheSigner(parsed, credentials, now) +} diff --git a/packages/daemon/test/source-cache-config.test.ts b/packages/daemon/test/source-cache-config.test.ts new file mode 100644 index 000000000..545a91474 --- /dev/null +++ b/packages/daemon/test/source-cache-config.test.ts @@ -0,0 +1,199 @@ +import { mkdtempSync, mkdirSync, writeFileSync } from 'node:fs' +import { tmpdir } from 'node:os' +import { join } from 'node:path' +import { afterEach, describe, expect, it } from 'vitest' +import { ConfigSchema } from '../src/config/config-schema.js' +import { loadConfig } from '../src/config/load-config.js' +import { + SOURCE_CACHE_ACCESS_KEY_ID_ENV, + SOURCE_CACHE_CONFIG_ENV, + SOURCE_CACHE_SECRET_ACCESS_KEY_ENV, + SOURCE_CACHE_SESSION_TOKEN_ENV +} from '../src/source-cache/config.js' + +const SOURCE_CACHE_ENV_NAMES = [ + SOURCE_CACHE_CONFIG_ENV, + SOURCE_CACHE_ACCESS_KEY_ID_ENV, + SOURCE_CACHE_SECRET_ACCESS_KEY_ENV, + SOURCE_CACHE_SESSION_TOKEN_ENV +] as const + +const saved = new Map() + +function withSourceCacheEnv( + values: Partial>, + run: () => void +): void { + for (const name of SOURCE_CACHE_ENV_NAMES) { + if (!saved.has(name)) saved.set(name, process.env[name]) + delete process.env[name] + } + for (const [name, value] of Object.entries(values)) { + if (value === undefined) delete process.env[name] + else process.env[name] = value + } + try { + run() + } finally { + for (const name of SOURCE_CACHE_ENV_NAMES) { + const value = saved.get(name) + if (value === undefined) delete process.env[name] + else process.env[name] = value + } + } +} + +function tmpRoot(): string { + const root = mkdtempSync(join(tmpdir(), 'ac-source-cache-cfg-')) + mkdirSync(root, { recursive: true }) + writeFileSync(join(root, 'config.json'), JSON.stringify({ version: 1 })) + return root +} + +afterEach(() => { + for (const [name, value] of saved) { + if (value === undefined) delete process.env[name] + else process.env[name] = value + } + saved.clear() +}) + +describe('source cache daemon configuration', () => { + it('is absent, and therefore a no-op, when no bucket is configured', () => { + withSourceCacheEnv({}, () => { + const cfg = loadConfig({ root: tmpRoot() }) + expect(cfg.sourceCache).toBeUndefined() + }) + }) + + it('validates the member-only environment document and resolves Secret credentials', () => { + withSourceCacheEnv( + { + [SOURCE_CACHE_CONFIG_ENV]: JSON.stringify({ + endpoint: 'https://cache.example.test', + region: 'us-east-1', + bucket: 'agentconnect-cache', + prefix: '/install-a/', + forcePathStyle: false, + credentialSource: 'secret', + limits: { + maxBundleBytes: 3_000, + orgTotalBytes: 30_000, + pendingReservationMs: 7_200_000, + pendingObjectTtlDays: 3, + unreferencedObjectTtlDays: 8, + unreadPointerTtlDays: 31 + } + }), + [SOURCE_CACHE_ACCESS_KEY_ID_ENV]: 'AKIDEXAMPLE', + [SOURCE_CACHE_SECRET_ACCESS_KEY_ENV]: 'secret-example', + [SOURCE_CACHE_SESSION_TOKEN_ENV]: 'session-example' + }, + () => { + const cfg = loadConfig({ root: tmpRoot() }) + expect(cfg.sourceCache).toEqual({ + endpoint: 'https://cache.example.test', + region: 'us-east-1', + bucket: 'agentconnect-cache', + prefix: 'install-a', + forcePathStyle: false, + credentials: { + source: 'secret', + accessKeyId: 'AKIDEXAMPLE', + secretAccessKey: 'secret-example', + sessionToken: 'session-example' + }, + limits: { + maxBundleBytes: 3_000, + orgTotalBytes: 30_000, + pendingReservationMs: 7_200_000, + pendingObjectTtlDays: 3, + unreferencedObjectTtlDays: 8, + unreadPointerTtlDays: 31 + } + }) + } + ) + }) + + it("applies the design's defaults, including the ServiceAccount credential source", () => { + withSourceCacheEnv( + { + [SOURCE_CACHE_CONFIG_ENV]: JSON.stringify({ + endpoint: 'https://cache.example.test', + bucket: 'agentconnect-cache', + credentialSource: 'serviceAccount' + }) + }, + () => { + const cfg = loadConfig({ root: tmpRoot() }) + expect(cfg.sourceCache?.credentials).toEqual({ source: 'serviceAccount' }) + expect(cfg.sourceCache?.region).toBe('auto') + expect(cfg.sourceCache?.prefix).toBe('') + expect(cfg.sourceCache?.forcePathStyle).toBe(false) + expect(cfg.sourceCache?.limits).toEqual({ + maxBundleBytes: 2 * 1024 ** 3, + orgTotalBytes: 20 * 1024 ** 3, + pendingReservationMs: 60 * 60 * 1000, + pendingObjectTtlDays: 2, + unreferencedObjectTtlDays: 7, + unreadPointerTtlDays: 30 + }) + } + ) + }) + + it('rejects malformed config rather than disabling silently', () => { + withSourceCacheEnv( + { + [SOURCE_CACHE_CONFIG_ENV]: JSON.stringify({ + endpoint: 'not a url', + bucket: 'agentconnect-cache', + credentialSource: 'serviceAccount' + }) + }, + () => { + expect(() => loadConfig({ root: tmpRoot() })).toThrow(SOURCE_CACHE_CONFIG_ENV) + } + ) + }) + + it('requires both static keys when the credential source is a Secret', () => { + withSourceCacheEnv( + { + [SOURCE_CACHE_CONFIG_ENV]: JSON.stringify({ + endpoint: 'https://cache.example.test', + bucket: 'agentconnect-cache', + credentialSource: 'secret' + }), + [SOURCE_CACHE_ACCESS_KEY_ID_ENV]: 'AKIDEXAMPLE' + }, + () => { + expect(() => loadConfig({ root: tmpRoot() })).toThrow(SOURCE_CACHE_SECRET_ACCESS_KEY_ENV) + } + ) + }) + + it('accepts a file document and keeps the environment out of the direct schema', () => { + expect(ConfigSchema.parse({ version: 1 })).not.toHaveProperty('sourceCache') + expect( + ConfigSchema.parse({ + version: 1, + sourceCache: { + endpoint: 'https://cache.example.test', + region: 'auto', + bucket: 'agentconnect-cache', + prefix: '', + forcePathStyle: true, + credentials: { source: 'serviceAccount' }, + limits: {} + } + }).sourceCache + ).toMatchObject({ + endpoint: 'https://cache.example.test', + forcePathStyle: true, + credentials: { source: 'serviceAccount' }, + limits: { maxBundleBytes: 2 * 1024 ** 3 } + }) + }) +}) diff --git a/packages/daemon/test/source-cache-minio.it.test.ts b/packages/daemon/test/source-cache-minio.it.test.ts new file mode 100644 index 000000000..a6443c362 --- /dev/null +++ b/packages/daemon/test/source-cache-minio.it.test.ts @@ -0,0 +1,63 @@ +import { createHash, randomUUID } from 'node:crypto' +import { describe, expect, it } from 'vitest' +import { SourceCacheConfigSchema } from '../src/source-cache/config.js' +import { createSourceCacheSigner } from '../src/source-cache/signer.js' + +const ENABLED = process.env.AC_SOURCE_CACHE_MINIO === '1' +const ENDPOINT = process.env.AC_SOURCE_CACHE_MINIO_ENDPOINT ?? 'http://127.0.0.1:19000' +const ACCESS_KEY_ID = process.env.AC_SOURCE_CACHE_MINIO_ACCESS_KEY ?? 'minioadmin' +const SECRET_ACCESS_KEY = process.env.AC_SOURCE_CACHE_MINIO_SECRET_KEY ?? 'minioadmin' +const BUCKET = process.env.AC_SOURCE_CACHE_MINIO_BUCKET ?? 'agentconnect-source-cache' + +describe.skipIf(!ENABLED)('source cache presigner against MinIO', () => { + it('accepts the signed PUT headers and serves the signed GET', async () => { + const signer = createSourceCacheSigner( + SourceCacheConfigSchema.parse({ + endpoint: ENDPOINT, + region: 'us-east-1', + bucket: BUCKET, + prefix: 'it', + forcePathStyle: true, + credentials: { + source: 'secret', + accessKeyId: ACCESS_KEY_ID, + secretAccessKey: SECRET_ACCESS_KEY + }, + limits: {} + }) + )! + const key = `src/contract/${randomUUID()}.bundle` + const body = Buffer.from(`source-cache-presign-${randomUUID()}`) + const checksumSha256 = createHash('sha256').update(body).digest('base64') + + const putUrl = await signer.presignPut(key, { + contentLength: body.byteLength, + checksumSha256 + }) + const tampered = await fetch(putUrl, { + method: 'PUT', + headers: { + 'content-length': String(body.byteLength), + 'x-amz-checksum-sha256': checksumSha256, + 'x-amz-tagging': 'ac-cache=live' + }, + body + }) + expect(tampered.status, await tampered.text()).toBe(403) + + const put = await fetch(putUrl, { + method: 'PUT', + headers: { + 'content-length': String(body.byteLength), + 'x-amz-checksum-sha256': checksumSha256, + 'x-amz-tagging': 'ac-cache=pending' + }, + body + }) + expect(put.status, await put.text()).toBe(200) + + const get = await fetch(await signer.presignGet(key)) + expect(get.status).toBe(200) + expect(Buffer.from(await get.arrayBuffer())).toEqual(body) + }) +}) diff --git a/packages/daemon/test/source-cache-signer.test.ts b/packages/daemon/test/source-cache-signer.test.ts new file mode 100644 index 000000000..6c7af3a83 --- /dev/null +++ b/packages/daemon/test/source-cache-signer.test.ts @@ -0,0 +1,126 @@ +import { describe, expect, it } from 'vitest' +import { + createSourceCacheSigner, + SOURCE_CACHE_GET_TTL_SECONDS, + SOURCE_CACHE_PUT_TTL_SECONDS +} from '../src/source-cache/signer.js' +import { SourceCacheConfigSchema, type SourceCacheConfig } from '../src/source-cache/config.js' + +const NOW = new Date('2026-10-01T15:04:05.000Z') + +function config(overrides: Partial = {}): SourceCacheConfig { + return SourceCacheConfigSchema.parse({ + endpoint: 'https://cache.example.test', + region: 'us-east-1', + bucket: 'agentconnect-cache', + prefix: 'install-a', + forcePathStyle: true, + credentials: { + source: 'secret', + accessKeyId: 'AKIDEXAMPLE', + secretAccessKey: 'wJalrXUtnFEMI/K7MDENG+bPxRfiCYEXAMPLEKEY' + }, + limits: {}, + ...overrides + }) +} + +describe('source cache presigner', () => { + it('returns no signer at all when the feature is disabled', () => { + expect(createSourceCacheSigner(undefined)).toBeUndefined() + }) + + it('presigns a GET for five minutes with only host in the signed headers', async () => { + const signer = createSourceCacheSigner(config(), { now: () => NOW })! + const url = new URL(await signer.presignGet('src/org/anon/repo/refs/hash/full/latest')) + + expect(url.origin).toBe('https://cache.example.test') + expect(url.pathname).toBe('/agentconnect-cache/install-a/src/org/anon/repo/refs/hash/full/latest') + expect(url.searchParams.get('X-Amz-Algorithm')).toBe('AWS4-HMAC-SHA256') + expect(url.searchParams.get('X-Amz-Credential')).toBe('AKIDEXAMPLE/20261001/us-east-1/s3/aws4_request') + expect(url.searchParams.get('X-Amz-Date')).toBe('20261001T150405Z') + expect(url.searchParams.get('X-Amz-Expires')).toBe(String(SOURCE_CACHE_GET_TTL_SECONDS)) + expect(url.searchParams.get('X-Amz-SignedHeaders')).toBe('host') + expect(url.searchParams.get('X-Amz-Signature')).toMatch(/^[0-9a-f]{64}$/) + }) + + it('presigns a PUT for fifteen minutes and binds length, checksum, and pending tagging', async () => { + const signer = createSourceCacheSigner(config(), { now: () => NOW })! + const url = new URL( + await signer.presignPut('src/org/anon/repo/bundles/00000000-0000-4000-8000-000000000001.bundle', { + contentLength: 12_345, + checksumSha256: '47DEQpj8HBSa+/TImW+5JCeuQeRkm5NMpJWZG3hSuFU=' + }) + ) + + expect(url.searchParams.get('X-Amz-Expires')).toBe(String(SOURCE_CACHE_PUT_TTL_SECONDS)) + expect(url.searchParams.get('X-Amz-SignedHeaders')).toBe('content-length;host;x-amz-checksum-sha256;x-amz-tagging') + }) + + it('encodes an object path exactly once', async () => { + const signer = createSourceCacheSigner(config(), { now: () => NOW })! + const url = new URL(await signer.presignGet('src/org/anon/repo/a file.bundle')) + + expect(url.pathname).toBe('/agentconnect-cache/install-a/src/org/anon/repo/a%20file.bundle') + expect(url.searchParams.get('X-Amz-Signature')).toMatch(/^[0-9a-f]{64}$/) + }) + + it('allows a caller to shorten a lifetime but never past the SigV4 seven-day ceiling', async () => { + const signer = createSourceCacheSigner(config(), { now: () => NOW })! + const short = new URL(await signer.presignGet('src/org/anon/repo/refs/hash/full/latest', { ttlSeconds: 60 })) + expect(short.searchParams.get('X-Amz-Expires')).toBe('60') + + await expect( + signer.presignPut('src/org/anon/repo/bundles/x.bundle', { + contentLength: 1, + checksumSha256: '47DEQpj8HBSa+/TImW+5JCeuQeRkm5NMpJWZG3hSuFU=', + ttlSeconds: 7 * 24 * 60 * 60 + 1 + }) + ).rejects.toThrow(/seven days|604800/i) + }) + + it('uses the Kubernetes service account credential chain without embedding credentials in config', async () => { + const signer = createSourceCacheSigner( + config({ + credentials: { source: 'serviceAccount' } + }), + { + now: () => NOW, + env: { + AWS_ACCESS_KEY_ID: 'SERVICEACCOUNT_KEY', + AWS_SECRET_ACCESS_KEY: 'SERVICEACCOUNT_SECRET', + AWS_SESSION_TOKEN: 'SERVICEACCOUNT_TOKEN' + } + } + )! + const url = new URL(await signer.presignGet('src/org/anon/repo/refs/hash/full/latest')) + + expect(url.searchParams.get('X-Amz-Credential')).toContain('SERVICEACCOUNT_KEY/') + expect(url.searchParams.get('X-Amz-Security-Token')).toBe('SERVICEACCOUNT_TOKEN') + }) + + it('uses path-style addressing for an IPv6 endpoint', async () => { + const signer = createSourceCacheSigner( + config({ + endpoint: 'http://[::1]:9000', + forcePathStyle: false + }), + { now: () => NOW } + )! + const url = new URL(await signer.presignGet('src/org/anon/repo/refs/hash/full/latest')) + + expect(url.host).toBe('[::1]:9000') + expect(url.pathname).toBe('/agentconnect-cache/install-a/src/org/anon/repo/refs/hash/full/latest') + }) + + it('rejects unsafe object keys and malformed checksums before signing', async () => { + const signer = createSourceCacheSigner(config(), { now: () => NOW })! + await expect(signer.presignGet('../escape')).rejects.toThrow(/object key/i) + await expect( + signer.presignPut('src/org/anon/repo/bundles/x.bundle', { + contentLength: 1, + checksumSha256: 'not-a-sha256' + }) + ).rejects.toThrow(/checksum/i) + }) +}) diff --git a/scripts/test-chart-render.rb b/scripts/test-chart-render.rb index 8485a4fc1..2f99c02b6 100644 --- a/scripts/test-chart-render.rb +++ b/scripts/test-chart-render.rb @@ -610,4 +610,94 @@ .dig('spec', 'template', 'spec', 'containers', 0).fetch('env').map { |item| item.fetch('name') } abort('an unset grace must leave both sides on the daemon default') if default_member.include?('AC_K8S_ORPHAN_GRACE_MS') +# Source Cache is a daemon-pool member capability only. The default has no bucket configured, so +# it must render no source-cache environment at all; enabling Secret credentials must project +# values by reference, never inline, and must stay out of reconciler/sandbox objects. +abort('source Cache must be off by default') if container.fetch('env').any? { |item| item.fetch('name').start_with?('AC_SOURCE_CACHE_') } + +source_cache_set = [ + '--set', 'sourceCache.endpoint=https://cache.example.test', + '--set', 'sourceCache.region=us-east-1', + '--set', 'sourceCache.bucket=agentconnect-cache', + '--set', 'sourceCache.prefix=install-a', + '--set', 'sourceCache.forcePathStyle=true', + '--set', 'sourceCache.credentialSource=secret', + '--set', 'sourceCache.existingSecret=source-cache-credentials', + '--set-json', 'sourceCache.limits={"maxBundleBytes":3000,"orgTotalBytes":30000,"pendingReservationMs":7200000,"pendingObjectTtlDays":3,"unreferencedObjectTtlDays":8,"unreadPointerTtlDays":31}' +] +source_rendered, source_error, source_status = Open3.capture3(*(command + source_cache_set)) +abort("helm template (source cache enabled) failed:\n#{source_error}") unless source_status.success? +source_documents = YAML.load_stream(source_rendered).compact +source_deployment = source_documents.find do |doc| + doc['kind'] == 'Deployment' && doc.dig('metadata', 'name') == 'example-agentconnect-daemon-pool' +end || abort('missing daemon-pool Deployment in the source-cache render') +source_container = source_deployment.dig('spec', 'template', 'spec', 'containers').find { |item| item['name'] == 'daemon-pool' } || + abort('missing daemon-pool container in the source-cache render') +source_env = source_container.fetch('env') +config_entry = source_env.find { |item| item['name'] == 'AC_SOURCE_CACHE_CONFIG' } || + abort('daemon pool must receive the source-cache config') +source_config = JSON.parse(config_entry.fetch('value')) +abort('source-cache config must carry the endpoint') unless source_config['endpoint'] == 'https://cache.example.test' +abort('source-cache config must carry the region') unless source_config['region'] == 'us-east-1' +abort('source-cache config must carry the bucket') unless source_config['bucket'] == 'agentconnect-cache' +abort('source-cache config must carry the prefix') unless source_config['prefix'] == 'install-a' +abort('source-cache config must carry forcePathStyle') unless source_config['forcePathStyle'] == true +abort('source-cache config must name the credential source') unless source_config['credentialSource'] == 'secret' +abort('source-cache config must carry section 10 limits') unless source_config['limits'] == { + 'maxBundleBytes' => 3000, + 'orgTotalBytes' => 30000, + 'pendingReservationMs' => 7_200_000, + 'pendingObjectTtlDays' => 3, + 'unreferencedObjectTtlDays' => 8, + 'unreadPointerTtlDays' => 31 +} + +{ + 'AC_SOURCE_CACHE_ACCESS_KEY_ID' => 'AWS_ACCESS_KEY_ID', + 'AC_SOURCE_CACHE_SECRET_ACCESS_KEY' => 'AWS_SECRET_ACCESS_KEY', + 'AC_SOURCE_CACHE_SESSION_TOKEN' => 'AWS_SESSION_TOKEN' +}.each do |env_name, secret_key| + entry = source_env.find { |item| item['name'] == env_name } || abort("daemon pool must project #{env_name}") + abort("#{env_name} must be a Secret reference, never an inline value") if entry.key?('value') + abort("#{env_name} must read #{secret_key} from sourceCache.existingSecret") unless + entry.dig('valueFrom', 'secretKeyRef') == { + 'name' => 'source-cache-credentials', + 'key' => secret_key + } +end + +source_reconciler = source_documents.find do |doc| + doc['kind'] == 'CronJob' && doc.dig('metadata', 'name') == 'example-agentconnect-daemon-pool-reconciler' +end || abort('missing daemon-pool reconciler') +abort('the reconciler must not receive source-cache credentials') if source_reconciler.to_yaml.include?('AC_SOURCE_CACHE') +source_runtime = source_documents.find do |doc| + doc['kind'] == 'SandboxTemplate' && doc.dig('metadata', 'name') == 'example-agentconnect-runtime' +end || abort('missing runtime SandboxTemplate') +abort('sandbox pods must never receive source-cache credentials') if source_runtime.to_yaml.include?('AC_SOURCE_CACHE') + +service_account_set = [ + '--set', 'sourceCache.endpoint=https://cache.example.test', + '--set', 'sourceCache.bucket=agentconnect-cache', + '--set', 'sourceCache.credentialSource=serviceAccount' +] +service_rendered, service_error, service_status = Open3.capture3(*(command + service_account_set)) +abort("helm template (source cache service account) failed:\n#{service_error}") unless service_status.success? +service_documents = YAML.load_stream(service_rendered).compact +service_member = service_documents.find do |doc| + doc['kind'] == 'Deployment' && doc.dig('metadata', 'name') == 'example-agentconnect-daemon-pool' +end.dig('spec', 'template', 'spec', 'containers').find { |item| item['name'] == 'daemon-pool' } +service_env_names = service_member.fetch('env').map { |item| item.fetch('name') } +abort('ServiceAccount source must not inject static-key environment') if service_env_names.any? { |name| name.start_with?('AC_SOURCE_CACHE_ACCESS_KEY', 'AC_SOURCE_CACHE_SECRET') } + +[ + [['sourceCache.endpoint=https://cache.example.test', 'sourceCache.bucket=agentconnect-cache', 'sourceCache.credentialSource=secret'], 'existingSecret'], + [['sourceCache.endpoint=https://cache.example.test'], 'bucket'] +].each do |settings, expected| + extra = settings.flat_map { |setting| ['--set', setting] } + _, refused, refused_status = Open3.capture3(*(command + extra)) + abort("incomplete source cache #{settings.join(', ')} must be refused") if refused_status.success? + abort("refusal for #{settings.join(', ')} must say what is missing:\n#{refused}") unless refused.include?(expected) +end + + puts 'chart render contract: ok' From 45d9a61ad586b655356ce44e70221cf064c8ca8d Mon Sep 17 00:00:00 2001 From: Harbor404 <2657212322@qq.com> Date: Thu, 1 Oct 2026 18:26:57 +0800 Subject: [PATCH 2/3] fix(source-cache): isolate pool-member configuration --- charts/agentconnect/README.md | 15 + .../agentconnect/templates/daemon-pool.yaml | 107 ++++-- charts/agentconnect/values.yaml | 21 +- packages/daemon/src/config/config-schema.ts | 4 - packages/daemon/src/config/load-config.ts | 8 - packages/daemon/src/daemon.ts | 12 + packages/daemon/src/source-cache/config.ts | 143 +++++--- packages/daemon/src/source-cache/signer.ts | 27 +- packages/daemon/test/daemon-k8s-mode.test.ts | 30 ++ .../daemon/test/source-cache-config.test.ts | 343 ++++++++++-------- .../daemon/test/source-cache-signer.test.ts | 19 +- scripts/test-chart-render.rb | 122 +++++-- 12 files changed, 533 insertions(+), 318 deletions(-) diff --git a/charts/agentconnect/README.md b/charts/agentconnect/README.md index b33ec21c1..9dffeb305 100644 --- a/charts/agentconnect/README.md +++ b/charts/agentconnect/README.md @@ -67,6 +67,21 @@ helm install agentconnect oci://ghcr.io/agentconnect-md/charts/agentconnect \ --set daemonPool.enabled=false --set installCRD=false --set relay.enabled=false ``` +## Source Cache + +Source Cache is optional and off by default. Set `sourceCache.endpoint` and +`sourceCache.bucket` to enable it; only daemon-pool members receive the rendered +document, and the reconciler and sandbox pods receive neither the document nor +credentials. + +`sourceCache.credentialSource=serviceAccount` supports AWS identities only: the +AWS access/secret environment pair, IRSA/web identity, or the ECS/EKS Pod Identity +container endpoint. It also requires a concrete AWS region. Stores such as R2, +MinIO, and GKE Workload Identity use `credentialSource=secret`; the referenced +Secret is mounted as member-local files and is never placed in the process +environment. `sourceCache.sessionTokenKey` is empty by default for the ordinary +two-key access/secret pair and is set only when that Secret carries a session token. + ## Requirements - **Kubernetes >= 1.28** (the relay reads the `apps.kubernetes.io/pod-index` label). diff --git a/charts/agentconnect/templates/daemon-pool.yaml b/charts/agentconnect/templates/daemon-pool.yaml index 3e92bd558..b887abb29 100644 --- a/charts/agentconnect/templates/daemon-pool.yaml +++ b/charts/agentconnect/templates/daemon-pool.yaml @@ -5,6 +5,8 @@ {{- $sourceCache := .Values.sourceCache }} {{- $sourceCacheEndpoint := trim ($sourceCache.endpoint | default "") }} {{- $sourceCacheBucket := trim ($sourceCache.bucket | default "") }} +{{- $sourceCacheRegion := trim ($sourceCache.region | default "us-east-1") }} +{{- $sourceCacheCredentialSource := $sourceCache.credentialSource | default "serviceAccount" }} {{- $sourceCacheEnabled := and (ne $sourceCacheEndpoint "") (ne $sourceCacheBucket "") }} {{- if and (ne $sourceCacheEndpoint "") (eq $sourceCacheBucket "") }} {{- fail "sourceCache.bucket is required when sourceCache.endpoint is set" }} @@ -13,24 +15,58 @@ {{- fail "sourceCache.endpoint is required when sourceCache.bucket is set" }} {{- end }} {{- if $sourceCacheEnabled }} -{{- if not (has $sourceCache.credentialSource (list "serviceAccount" "secret")) }} +{{- if eq $sourceCacheRegion "" }} +{{- fail "sourceCache.region is required when source cache is enabled" }} +{{- end }} +{{- if not (has $sourceCacheCredentialSource (list "serviceAccount" "secret")) }} {{- fail "sourceCache.credentialSource must be serviceAccount or secret" }} {{- end }} -{{- if eq $sourceCache.credentialSource "secret" }} +{{- if eq $sourceCacheCredentialSource "secret" }} {{- if not $sourceCache.existingSecret }} {{- fail "sourceCache.existingSecret is required when sourceCache.credentialSource=secret" }} {{- end }} {{- if or (not $sourceCache.accessKeyIdKey) (not $sourceCache.secretAccessKeyKey) }} {{- fail "sourceCache.accessKeyIdKey and sourceCache.secretAccessKeyKey are required for Secret credentials" }} {{- end }} -{{- else if $sourceCache.existingSecret }} +{{- else }} +{{- if $sourceCache.existingSecret }} {{- fail "sourceCache.existingSecret is only valid when sourceCache.credentialSource=secret" }} {{- end }} -{{- range $name := list "AC_SOURCE_CACHE_CONFIG" "AC_SOURCE_CACHE_ACCESS_KEY_ID" "AC_SOURCE_CACHE_SECRET_ACCESS_KEY" "AC_SOURCE_CACHE_SESSION_TOKEN" }} -{{- if hasKey $.Values.daemonPool.extraEnv $name }} -{{- fail (printf "daemonPool.extraEnv.%s collides with the chart-owned Source Cache configuration" $name) }} +{{- if eq $sourceCacheRegion "auto" }} +{{- fail "sourceCache.region must be a concrete AWS region when sourceCache.credentialSource=serviceAccount; use secret for stores such as R2" }} +{{- end }} {{- end }} {{- end }} +{{- if $sourceCacheEnabled }} +{{- $sourceCacheCredentials := dict "source" $sourceCacheCredentialSource }} +{{- if eq $sourceCacheCredentialSource "secret" }} +{{- $_ := set $sourceCacheCredentials "accessKeyIdFile" "/var/run/ac-source-cache-credentials/access-key-id" }} +{{- $_ := set $sourceCacheCredentials "secretAccessKeyFile" "/var/run/ac-source-cache-credentials/secret-access-key" }} +{{- if $sourceCache.sessionTokenKey }} +{{- $_ := set $sourceCacheCredentials "sessionTokenFile" "/var/run/ac-source-cache-credentials/session-token" }} +{{- end }} +{{- end }} +{{- $sourceCacheDocument := dict + "version" 1 + "endpoint" $sourceCacheEndpoint + "region" $sourceCacheRegion + "bucket" $sourceCacheBucket + "prefix" ($sourceCache.prefix | default "") + "forcePathStyle" ($sourceCache.forcePathStyle | default false) + "credentials" $sourceCacheCredentials + "limits" $sourceCache.limits }} +# The member-only Source Cache document. It contains no static credential values; the optional +# credential Secret below is mounted as files and is never projected into the process environment. +apiVersion: v1 +kind: ConfigMap +metadata: + name: {{ include "agentconnect.fullname" . }}-source-cache + labels: + {{- include "agentconnect.labels" . | nindent 4 }} + {{- include "agentconnect.componentLabels" (dict "ctx" . "component" "daemon-pool") | nindent 4 }} +data: + config.json: {{ $sourceCacheDocument | toJson | quote }} +--- {{- end }} # The install-wide daemon pool. Its Pods are pool members. apiVersion: v1 @@ -195,31 +231,6 @@ spec: {{- include "agentconnect.daemonPoolClusterEnv" . | nindent 12 }} - name: AC_K8S_SHIM_PORT value: {{ .Values.daemonPool.shimPort | quote }} - {{- if $sourceCacheEnabled }} - # Member-only Source Cache configuration. Credential values are always projected - # by Secret reference; the static key path below never renders a value inline. - - name: AC_SOURCE_CACHE_CONFIG - value: {{ dict "endpoint" $sourceCacheEndpoint "region" $sourceCache.region "bucket" $sourceCacheBucket "prefix" $sourceCache.prefix "forcePathStyle" $sourceCache.forcePathStyle "credentialSource" $sourceCache.credentialSource "limits" $sourceCache.limits | toJson | quote }} - {{- if eq $sourceCache.credentialSource "secret" }} - - name: AC_SOURCE_CACHE_ACCESS_KEY_ID - valueFrom: - secretKeyRef: - name: {{ $sourceCache.existingSecret | quote }} - key: {{ $sourceCache.accessKeyIdKey | quote }} - - name: AC_SOURCE_CACHE_SECRET_ACCESS_KEY - valueFrom: - secretKeyRef: - name: {{ $sourceCache.existingSecret | quote }} - key: {{ $sourceCache.secretAccessKeyKey | quote }} - {{- if $sourceCache.sessionTokenKey }} - - name: AC_SOURCE_CACHE_SESSION_TOKEN - valueFrom: - secretKeyRef: - name: {{ $sourceCache.existingSecret | quote }} - key: {{ $sourceCache.sessionTokenKey | quote }} - {{- end }} - {{- end }} - {{- end }} # The sink the readiness probe below reads: the daemon serves GET /readyz here and # answers 200 only once startup finished, the control plane acknowledged registration # and the install-wide runtime probe returned — 503 otherwise, including from SIGTERM @@ -414,6 +425,16 @@ spec: - name: data-plane mountPath: /var/run/ac-data-plane readOnly: true + {{- if $sourceCacheEnabled }} + - name: source-cache + mountPath: /var/run/ac-source-cache + readOnly: true + {{- if eq $sourceCacheCredentialSource "secret" }} + - name: source-cache-credentials + mountPath: /var/run/ac-source-cache-credentials + readOnly: true + {{- end }} + {{- end }} - name: cp-identity mountPath: /var/run/ac-cp-identity readOnly: true @@ -439,6 +460,30 @@ spec: items: - key: {{ .Values.daemonPool.dataPlane.key }} path: config.json + {{- if $sourceCacheEnabled }} + - name: source-cache + configMap: + name: {{ include "agentconnect.fullname" . }}-source-cache + defaultMode: 0400 + items: + - key: config.json + path: config.json + {{- if eq $sourceCacheCredentialSource "secret" }} + - name: source-cache-credentials + secret: + secretName: {{ $sourceCache.existingSecret | quote }} + defaultMode: 0400 + items: + - key: {{ $sourceCache.accessKeyIdKey | quote }} + path: access-key-id + - key: {{ $sourceCache.secretAccessKeyKey | quote }} + path: secret-access-key + {{- if $sourceCache.sessionTokenKey }} + - key: {{ $sourceCache.sessionTokenKey | quote }} + path: session-token + {{- end }} + {{- end }} + {{- end }} - name: cp-identity projected: defaultMode: 0400 diff --git a/charts/agentconnect/values.yaml b/charts/agentconnect/values.yaml index 13be27cca..e4b719084 100644 --- a/charts/agentconnect/values.yaml +++ b/charts/agentconnect/values.yaml @@ -410,27 +410,32 @@ daemonPool: tolerations: [] # Optional Source Cache. Empty endpoint or bucket leaves it OFF: no URLs are issued and -# pool members clone from origin exactly as before. Only the daemon-pool member receives -# this configuration; sandbox pods never receive bucket credentials. +# pool members clone from origin exactly as before. Only the daemon-pool member receives the +# mounted document; sandbox pods and the reconciler receive neither the document nor credentials. sourceCache: # Full S3 API origin, e.g. https://s3.us-east-1.amazonaws.com or https://minio.example.test endpoint: '' - # SigV4 region. Use "auto" for R2; a concrete region for S3 or MinIO. - region: auto + # SigV4 region. serviceAccount requires a concrete AWS region. Use "auto" only with secret + # credentials for a non-AWS S3-compatible store such as R2. + region: us-east-1 # Bucket (or bucket prefix) dedicated to this install's Source Cache. bucket: '' # Optional key prefix inside the bucket. The daemon adds src/, files/, and snapshots/ below it. prefix: '' # MinIO and other S3-compatible stores commonly require path-style addressing. forcePathStyle: false - # serviceAccount uses the pool member's IAM/Workload Identity. secret projects static keys - # from an existing Secret by name; it never embeds credential values in this values file. + # serviceAccount resolves AWS identities only: the AWS access/secret environment pair, + # IRSA/web identity, or the ECS/EKS Pod Identity container endpoint. GKE Workload Identity, + # R2, MinIO, and other non-AWS stores require secret credentials. + # secret mounts an existing Kubernetes Secret as member-local files; credential values are + # never rendered into this values file, the Source Cache document, or process environment. credentialSource: serviceAccount existingSecret: '' accessKeyIdKey: AWS_ACCESS_KEY_ID secretAccessKeyKey: AWS_SECRET_ACCESS_KEY - # Optional; set to "" for stores or identities that do not issue a session token. - sessionTokenKey: AWS_SESSION_TOKEN + # Optional session-token file. Empty is the normal two-key access/secret setup; set it only + # when the referenced Secret actually carries a session token. + sessionTokenKey: '' # Section 10 limits. The daemon validates the values; the data plane enforces them. limits: maxBundleBytes: 2147483648 # 2 GiB diff --git a/packages/daemon/src/config/config-schema.ts b/packages/daemon/src/config/config-schema.ts index 2f1c48162..565fff55f 100644 --- a/packages/daemon/src/config/config-schema.ts +++ b/packages/daemon/src/config/config-schema.ts @@ -6,7 +6,6 @@ import { SESSION_RETENTION_RE, type SessionRetentionSetting } from '@agentconnect.md/protocol' -import { SourceCacheConfigSchema } from '../source-cache/config.js' /** The `{name, value}[]` shape shared by runtime env, MCP env, and MCP headers. */ const NameValueList = z.array(z.object({ name: z.string(), value: z.string() })).default([]) @@ -247,9 +246,6 @@ export const ConfigSchema = z.object({ // degradation); the CP's register/ok snapshot re-converges it authoritatively once // connected. CP-owned — overwritten on every roster converge, not hand-edited. relays: z.array(RelayRosterEntry).default([]), - // Optional Source Cache. Absence is the disabled no-op; the environment document is - // parsed by loadConfig and never exposed to sandboxes or the control plane. - sourceCache: SourceCacheConfigSchema.optional(), logging: z .object({ level: z.enum(['trace', 'debug', 'info', 'warn', 'error']).default('info') }) .default({ level: 'info' }), diff --git a/packages/daemon/src/config/load-config.ts b/packages/daemon/src/config/load-config.ts index 9b600814e..6aac7de72 100644 --- a/packages/daemon/src/config/load-config.ts +++ b/packages/daemon/src/config/load-config.ts @@ -8,7 +8,6 @@ import { } from '@agentconnect.md/protocol' import { z } from 'zod' import { ConfigSchema, WorkspaceGitOrigin, type Config } from './config-schema.js' -import { sourceCacheConfigFromEnv } from '../source-cache/config.js' import { resolveRoot, configPath, defaultAgentsDir } from '../paths.js' export interface FlatOverrides { @@ -94,13 +93,6 @@ export function loadConfig( throw new Error(`config not found: ${file} (create it, pass --config, or run \`agentconnect login\`)`) } const cfg = ConfigSchema.parse(raw) // throws on invalid - // Pool members have no writable config file. The chart supplies the Source Cache - // document by environment; a file that states sourceCache wins (same precedence as - // the other operator-owned policy overrides below). - if ((raw as { sourceCache?: unknown } | null)?.sourceCache === undefined) { - const sourceCache = sourceCacheConfigFromEnv(process.env) - if (sourceCache) cfg.sourceCache = sourceCache - } for (const warning of mapLegacySandbox(cfg)) opts.warn?.(warning) const o = opts.overrides ?? {} diff --git a/packages/daemon/src/daemon.ts b/packages/daemon/src/daemon.ts index 9e1d92ef1..75ca04aad 100644 --- a/packages/daemon/src/daemon.ts +++ b/packages/daemon/src/daemon.ts @@ -698,6 +698,8 @@ import { managedDistillCapture, withManagedDistill } from './memory/managed-dist import { defaultMemoryPluginMetrics } from './memory-plugin/metrics.js' import { openPostgresDataPlane, type PostgresDataPlane } from './store/postgres-data-plane.js' import { DATA_PLANE_CONFIG_PATH } from './store/postgres-config.js' +import { readSourceCacheConfig } from './source-cache/config.js' +import { createSourceCacheSigner, type SourceCacheSigner } from './source-cache/signer.js' import type { EvaluationCapabilityProfile } from './evaluation/events.js' import { DaemonEvaluationHooks, type DaemonEvaluationHost } from './evaluation/daemon-hooks.js' import { SessionMetadataOutbox, type SessionMetadataHost } from './store/session-metadata-outbox.js' @@ -1583,6 +1585,8 @@ export class Daemon { private readonly clusterIdentityToken?: () => string | undefined // The k8s execution plane: shim dialer + driver + workspace seam. Undefined outside --k8s. private k8sPlane?: K8sRuntimePlane + // The pool-member Source Cache signer. Undefined is the off/no-op state; only --k8s composes it. + private sourceCache?: SourceCacheSigner // An undefined entry awaits retry; a promise is the one takeover already in flight. private readonly k8sAdoptions = new Map | undefined>() private microsandbox?: MicrosandboxManager @@ -1852,6 +1856,8 @@ export class Daemon { startK8sPlane?: typeof startK8sRuntimePlane /** Test seam only; production reads the fixed Secret mount under `--k8s`, else the `postgres` store's file. */ openDataPlane?: typeof openPostgresDataPlane + /** Test seam only; production reads the fixed member-only Source Cache mount under `--k8s`. */ + readSourceCacheConfig?: typeof readSourceCacheConfig /** Test seam for the pool member startup barrier; production waits for CP register/ok. */ startControlPlane?: (root: string) => Promise | undefined /** Test seams for local catalog resolution and executable/state filtering. */ @@ -2702,6 +2708,12 @@ export class Daemon { /** Phase 4 — the shared data plane (under --k8s, or a `postgres` store), then under --k8s the execution plane the workspaces resolve through. */ private async startClusterPlanes(root: string, cfg: Config): Promise { + // Source Cache is a pool-member capability. Composing it here keeps the shared Config loader and + // every non-k8s command unaware of both the document and any static credentials it resolves. + if (this.k8s) { + this.sourceCache = createSourceCacheSigner((this.opts.readSourceCacheConfig ?? readSourceCacheConfig)()) + if (this.sourceCache) this.log.info('source cache: enabled') + } // `--k8s` needs the pool's shared store whatever the file says; any other daemon opens one only when its owner asked (#2188). const dataPlaneConfig = this.k8s ? DATA_PLANE_CONFIG_PATH diff --git a/packages/daemon/src/source-cache/config.ts b/packages/daemon/src/source-cache/config.ts index 9924011da..7b83312af 100644 --- a/packages/daemon/src/source-cache/config.ts +++ b/packages/daemon/src/source-cache/config.ts @@ -1,11 +1,12 @@ +import { readFileSync } from 'node:fs' import { z } from 'zod' -/** The one environment document the Helm chart renders into a pool member. */ -export const SOURCE_CACHE_CONFIG_ENV = 'AC_SOURCE_CACHE_CONFIG' -/** Static S3 credentials are projected by Secret reference, never inline in the config document. */ -export const SOURCE_CACHE_ACCESS_KEY_ID_ENV = 'AC_SOURCE_CACHE_ACCESS_KEY_ID' -export const SOURCE_CACHE_SECRET_ACCESS_KEY_ENV = 'AC_SOURCE_CACHE_SECRET_ACCESS_KEY' -export const SOURCE_CACHE_SESSION_TOKEN_ENV = 'AC_SOURCE_CACHE_SESSION_TOKEN' +/** The pool member's fixed Secret/ConfigMap mount. It is never read by the shared daemon Config loader. */ +export const SOURCE_CACHE_CONFIG_PATH = '/var/run/ac-source-cache/config.json' +/** Secret-volume paths for static credentials. Keeping these files member-local avoids process-env inheritance. */ +export const SOURCE_CACHE_ACCESS_KEY_ID_FILE = '/var/run/ac-source-cache-credentials/access-key-id' +export const SOURCE_CACHE_SECRET_ACCESS_KEY_FILE = '/var/run/ac-source-cache-credentials/secret-access-key' +export const SOURCE_CACHE_SESSION_TOKEN_FILE = '/var/run/ac-source-cache-credentials/session-token' const GIB = 1024 ** 3 const HOUR_MS = 60 * 60 * 1000 @@ -101,63 +102,115 @@ export const SourceCacheCredentialsSchema = z.discriminatedUnion('source', [ SecretCredentials ]) -const SourceCacheBaseSchema = z +const SourceCacheBaseShape = { + endpoint: Endpoint, + // `auto` is an R2-style value and has no AWS STS region meaning. Keep the default concrete so + // the default ServiceAccount arm remains a valid AWS identity configuration. + region: z.string().trim().min(1).max(64).default('us-east-1'), + bucket: BucketName, + prefix: Prefix.default(''), + forcePathStyle: z.boolean().default(false), + limits: SourceCacheLimitsSchema.default(SOURCE_CACHE_LIMIT_DEFAULTS) +} as const + +const sourceCacheRegionRefinement = ( + value: { region: string; credentials: { source: 'serviceAccount' | 'secret' } }, + ctx: z.RefinementCtx +): void => { + if (value.credentials.source === 'serviceAccount' && value.region === 'auto') { + ctx.addIssue({ + code: 'custom', + path: ['region'], + message: 'serviceAccount credentials require a concrete AWS region; region "auto" is only valid with secret credentials' + }) + } +} + +/** The resolved member shape. Presence means the feature is enabled; absence is the no-op. */ +export const SourceCacheConfigSchema = z .object({ - endpoint: Endpoint, - region: z.string().trim().min(1).max(64).default('auto'), - bucket: BucketName, - prefix: Prefix.default(''), - forcePathStyle: z.boolean().default(false), - limits: SourceCacheLimitsSchema.default(SOURCE_CACHE_LIMIT_DEFAULTS) + ...SourceCacheBaseShape, + credentials: SourceCacheCredentialsSchema }) .strict() - -/** The daemon's validated shape. Presence means the feature is enabled; absence is the no-op. */ -export const SourceCacheConfigSchema = SourceCacheBaseSchema.extend({ - credentials: SourceCacheCredentialsSchema -}).strict() + .superRefine(sourceCacheRegionRefinement) export type SourceCacheConfig = z.infer export type SourceCacheCredentials = SourceCacheConfig['credentials'] -const SourceCacheEnvironmentSchema = SourceCacheBaseSchema.extend({ - credentialSource: z.enum(['serviceAccount', 'secret']) -}).strict() +const CredentialFilePath = z + .string() + .min(1) + .refine((value) => value.startsWith('/'), 'credential file paths must be absolute') + +const SecretFileCredentials = z + .object({ + source: z.literal('secret'), + accessKeyIdFile: CredentialFilePath, + secretAccessKeyFile: CredentialFilePath, + sessionTokenFile: CredentialFilePath.optional() + }) + .strict() -function requiredSecret(env: NodeJS.ProcessEnv, name: string): string { - const value = env[name]?.trim() - if (!value) throw new Error(`${name} is required when source cache credentialSource=secret`) - return value +export const SourceCacheDocumentSchema = z + .object({ + version: z.literal(1), + ...SourceCacheBaseShape, + credentials: z.discriminatedUnion('source', [ServiceAccountCredentials, SecretFileCredentials]) + }) + .strict() + .superRefine(sourceCacheRegionRefinement) + +export type SourceCacheDocument = z.infer + +function readCredentialFile(path: string, label: string): string { + try { + const value = readFileSync(path, 'utf8').trim() + if (!value) throw new Error('empty') + return value + } catch { + throw new Error(`source cache ${label} file is not readable at ${path}`) + } } -/** Parse the member-only Helm document. Blank means the feature is off; malformed means - * configuration is wrong and must be fixed rather than silently disabling the cache. */ -export function sourceCacheConfigFromEnv(env: NodeJS.ProcessEnv = process.env): SourceCacheConfig | undefined { - const raw = env[SOURCE_CACHE_CONFIG_ENV]?.trim() - if (!raw) return undefined +/** Read the member-only document. A missing mount is the disabled no-op; a present malformed + * document fails startup instead of silently issuing unsigned or differently-scoped URLs. */ +export function readSourceCacheConfig(path = SOURCE_CACHE_CONFIG_PATH): SourceCacheConfig | undefined { + let text: string + try { + text = readFileSync(path, 'utf8') + } catch (error) { + if ((error as NodeJS.ErrnoException).code === 'ENOENT') return undefined + throw new Error(`source cache configuration is not readable at ${path}`) + } - let parsed: unknown + let raw: unknown try { - parsed = JSON.parse(raw) + raw = JSON.parse(text) } catch { - throw new Error(`${SOURCE_CACHE_CONFIG_ENV} must be valid JSON`) + throw new Error(`source cache configuration is not valid JSON at ${path}`) } - const document = SourceCacheEnvironmentSchema.safeParse(parsed) - if (!document.success) throw new Error(`${SOURCE_CACHE_CONFIG_ENV}: ${document.error.message}`) + const parsed = SourceCacheDocumentSchema.safeParse(raw) + if (!parsed.success) { + const issue = parsed.error.issues[0] + throw new Error( + `invalid source cache configuration at ${path}: ${issue?.path.join('.') || 'document'} ${issue?.message}` + ) + } - const credentials = - document.data.credentialSource === 'secret' + const { version: _version, credentials, ...base } = parsed.data + const resolvedCredentials: SourceCacheCredentials = + credentials.source === 'secret' ? { - source: 'secret' as const, - accessKeyId: requiredSecret(env, SOURCE_CACHE_ACCESS_KEY_ID_ENV), - secretAccessKey: requiredSecret(env, SOURCE_CACHE_SECRET_ACCESS_KEY_ENV), - ...(env[SOURCE_CACHE_SESSION_TOKEN_ENV]?.trim() - ? { sessionToken: env[SOURCE_CACHE_SESSION_TOKEN_ENV]!.trim() } + source: 'secret', + accessKeyId: readCredentialFile(credentials.accessKeyIdFile, 'access key id'), + secretAccessKey: readCredentialFile(credentials.secretAccessKeyFile, 'secret access key'), + ...(credentials.sessionTokenFile + ? { sessionToken: readCredentialFile(credentials.sessionTokenFile, 'session token') } : {}) } - : { source: 'serviceAccount' as const } + : { source: 'serviceAccount' } - const { credentialSource: _credentialSource, ...base } = document.data - return SourceCacheConfigSchema.parse({ ...base, credentials }) + return SourceCacheConfigSchema.parse({ ...base, credentials: resolvedCredentials }) } diff --git a/packages/daemon/src/source-cache/signer.ts b/packages/daemon/src/source-cache/signer.ts index 412db2a5a..5f0419df5 100644 --- a/packages/daemon/src/source-cache/signer.ts +++ b/packages/daemon/src/source-cache/signer.ts @@ -10,25 +10,19 @@ import { export const SOURCE_CACHE_GET_TTL_SECONDS = 5 * 60 export const SOURCE_CACHE_PUT_TTL_SECONDS = 15 * 60 -export const SOURCE_CACHE_MAX_PRESIGN_SECONDS = 7 * 24 * 60 * 60 - const UNSIGNED_PAYLOAD = 'UNSIGNED-PAYLOAD' const ALGORITHM = 'AWS4-HMAC-SHA256' const SERVICE = 's3' const TERMINATOR = 'aws4_request' const PENDING_TAG = 'ac-cache=pending' -export interface PresignOptions { - ttlSeconds?: number -} - -export interface PresignPutInput extends PresignOptions { +export interface PresignPutInput { contentLength: number checksumSha256: string } export interface SourceCacheSigner { - presignGet(key: string, options?: PresignOptions): Promise + presignGet(key: string): Promise presignPut(key: string, input: PresignPutInput): Promise } @@ -73,17 +67,6 @@ function formatAmzDate(date: Date): { timestamp: string; day: string } { return { timestamp, day: timestamp.slice(0, 8) } } -function ttlSeconds(value: number | undefined, fallback: number): number { - const ttl = value ?? fallback - if (!Number.isSafeInteger(ttl) || ttl <= 0) throw new Error('source cache presign TTL must be a positive integer') - if (ttl > SOURCE_CACHE_MAX_PRESIGN_SECONDS) { - throw new Error( - `source cache presign TTL cannot exceed the SigV4 seven-day ceiling (${SOURCE_CACHE_MAX_PRESIGN_SECONDS}s)` - ) - } - return ttl -} - function objectKey(prefix: string, key: string): string { if (!key || key !== key.trim() || key.startsWith('/') || key.endsWith('/')) throw new Error('invalid source cache object key') @@ -129,8 +112,8 @@ class AwsSigV4SourceCacheSigner implements SourceCacheSigner { this.endpoint = new URL(config.endpoint) } - async presignGet(key: string, options: PresignOptions = {}): Promise { - return this.presign('GET', key, {}, ttlSeconds(options.ttlSeconds, SOURCE_CACHE_GET_TTL_SECONDS)) + async presignGet(key: string): Promise { + return this.presign('GET', key, {}, SOURCE_CACHE_GET_TTL_SECONDS) } async presignPut(key: string, input: PresignPutInput): Promise { @@ -142,7 +125,7 @@ class AwsSigV4SourceCacheSigner implements SourceCacheSigner { 'x-amz-checksum-sha256': checksumSha256(input.checksumSha256), 'x-amz-tagging': PENDING_TAG }, - ttlSeconds(input.ttlSeconds, SOURCE_CACHE_PUT_TTL_SECONDS) + SOURCE_CACHE_PUT_TTL_SECONDS ) } diff --git a/packages/daemon/test/daemon-k8s-mode.test.ts b/packages/daemon/test/daemon-k8s-mode.test.ts index dc6ca2091..5cf040ede 100644 --- a/packages/daemon/test/daemon-k8s-mode.test.ts +++ b/packages/daemon/test/daemon-k8s-mode.test.ts @@ -19,6 +19,7 @@ import { DEFAULT_SHIM_WORKSPACE_ROOT } from '../src/shim/protocol.js' import type { ResolvedRuntimeCatalog } from '../src/runtimes/registry.js' import { LocalStore } from '../src/store/local-store.js' import { DATA_PLANE_CONFIG_PATH } from '../src/store/postgres-config.js' +import { SourceCacheConfigSchema, type readSourceCacheConfig } from '../src/source-cache/config.js' import { mcpSocketPath, statePath } from '../src/paths.js' import { SANDBOX_TUNNEL_PATHS } from '../src/shim/sandbox-paths.js' @@ -100,6 +101,7 @@ function daemon(opts: { /** A store shared by several members, so a pool-wide probe can be asserted across them. */ store?: LocalStore openDataPlane?: ReturnType + readSourceCacheConfig?: typeof readSourceCacheConfig startControlPlane?: ReturnType /** Receives the options the mode hands the plane — the rows about what it asks the pod to serve. */ onPlaneStart?: (options: any) => void @@ -133,6 +135,7 @@ function daemon(opts: { } : {}), ...(opts.openDataPlane ? { openDataPlane: opts.openDataPlane as never } : {}), + ...(opts.readSourceCacheConfig ? { readSourceCacheConfig: opts.readSourceCacheConfig as never } : {}), startControlPlane: (opts.startControlPlane ?? vi.fn(() => Promise.resolve())) as never, ...(opts.supervisor ? { supervisor: opts.supervisor } : {}), resolveCatalog: async () => catalog(), @@ -281,6 +284,33 @@ describe('daemon --k8s mode', () => { } }) + it('loads the source-cache document only on a k8s pool member', async () => { + const sourceCache = SourceCacheConfigSchema.parse({ + endpoint: 'https://cache.example.test', + region: 'us-east-1', + bucket: 'agentconnect-cache', + credentials: { source: 'serviceAccount' }, + limits: {} + }) + const readSourceCacheConfig = vi.fn(() => sourceCache) + + const cfg = { store: { backend: 'local' }, limits: { poolShutdownDrainMs: 1 } } as never + + const local = daemon({ root: root(), k8s: false, readSourceCacheConfig }) + await (local as any).startClusterPlanes(root(), cfg) + expect(readSourceCacheConfig).not.toHaveBeenCalled() + expect((local as any).sourceCache).toBeUndefined() + + const member = daemon({ root: root(), k8s: true, readSourceCacheConfig }) + try { + await (member as any).startClusterPlanes(root(), cfg) + expect(readSourceCacheConfig).toHaveBeenCalledTimes(1) + expect((member as any).sourceCache).toBeDefined() + } finally { + await (member as any).dataPlane?.close() + } + }) + it('opens a postgres store named by its file outside k8s mode, and keeps no SQLite database (#2188)', async () => { const rootDir = root({ store: { backend: 'postgres', configFile: 'data-plane.json' }, cp: true }) const openDataPlane = vi.fn(() => fakeDataPlane()) diff --git a/packages/daemon/test/source-cache-config.test.ts b/packages/daemon/test/source-cache-config.test.ts index 545a91474..aca3fbb32 100644 --- a/packages/daemon/test/source-cache-config.test.ts +++ b/packages/daemon/test/source-cache-config.test.ts @@ -1,199 +1,230 @@ import { mkdtempSync, mkdirSync, writeFileSync } from 'node:fs' import { tmpdir } from 'node:os' import { join } from 'node:path' -import { afterEach, describe, expect, it } from 'vitest' +import { afterEach, describe, expect, it, vi } from 'vitest' import { ConfigSchema } from '../src/config/config-schema.js' import { loadConfig } from '../src/config/load-config.js' import { - SOURCE_CACHE_ACCESS_KEY_ID_ENV, - SOURCE_CACHE_CONFIG_ENV, - SOURCE_CACHE_SECRET_ACCESS_KEY_ENV, - SOURCE_CACHE_SESSION_TOKEN_ENV + SOURCE_CACHE_ACCESS_KEY_ID_FILE, + SOURCE_CACHE_CONFIG_PATH, + SOURCE_CACHE_SECRET_ACCESS_KEY_FILE, + SourceCacheConfigSchema, + SourceCacheDocumentSchema, + readSourceCacheConfig } from '../src/source-cache/config.js' -const SOURCE_CACHE_ENV_NAMES = [ - SOURCE_CACHE_CONFIG_ENV, - SOURCE_CACHE_ACCESS_KEY_ID_ENV, - SOURCE_CACHE_SECRET_ACCESS_KEY_ENV, - SOURCE_CACHE_SESSION_TOKEN_ENV -] as const - -const saved = new Map() - -function withSourceCacheEnv( - values: Partial>, - run: () => void -): void { - for (const name of SOURCE_CACHE_ENV_NAMES) { - if (!saved.has(name)) saved.set(name, process.env[name]) - delete process.env[name] - } - for (const [name, value] of Object.entries(values)) { - if (value === undefined) delete process.env[name] - else process.env[name] = value - } - try { - run() - } finally { - for (const name of SOURCE_CACHE_ENV_NAMES) { - const value = saved.get(name) - if (value === undefined) delete process.env[name] - else process.env[name] = value - } - } -} - -function tmpRoot(): string { +function tmpRoot(extra: Record = {}): string { const root = mkdtempSync(join(tmpdir(), 'ac-source-cache-cfg-')) mkdirSync(root, { recursive: true }) - writeFileSync(join(root, 'config.json'), JSON.stringify({ version: 1 })) + writeFileSync(join(root, 'config.json'), JSON.stringify({ version: 1, ...extra })) return root } -afterEach(() => { - for (const [name, value] of saved) { - if (value === undefined) delete process.env[name] - else process.env[name] = value +function writeDocument(document: unknown, path = join(tmpRoot(), 'source-cache.json')): string { + writeFileSync(path, JSON.stringify(document)) + return path +} + +function writeSecretFiles(root: string, values: { accessKeyId?: string; secretAccessKey?: string; sessionToken?: string } = {}) { + const dir = join(root, 'credentials') + mkdirSync(dir, { recursive: true }) + const files = { + accessKeyIdFile: join(dir, 'access-key-id'), + secretAccessKeyFile: join(dir, 'secret-access-key'), + sessionTokenFile: join(dir, 'session-token') } - saved.clear() + writeFileSync(files.accessKeyIdFile, values.accessKeyId ?? 'AKIDEXAMPLE') + writeFileSync(files.secretAccessKeyFile, values.secretAccessKey ?? 'secret-example') + if (values.sessionToken === undefined) { + const { sessionTokenFile: _sessionTokenFile, ...withoutSessionToken } = files + return withoutSessionToken + } + writeFileSync(files.sessionTokenFile, values.sessionToken) + return files +} + +afterEach(() => { + vi.unstubAllEnvs() }) -describe('source cache daemon configuration', () => { - it('is absent, and therefore a no-op, when no bucket is configured', () => { - withSourceCacheEnv({}, () => { - const cfg = loadConfig({ root: tmpRoot() }) - expect(cfg.sourceCache).toBeUndefined() - }) +describe('source cache member configuration', () => { + it('is absent, and therefore a no-op, when the pool member has no document', () => { + expect(readSourceCacheConfig(join(tmpRoot(), 'missing.json'))).toBeUndefined() }) - it('validates the member-only environment document and resolves Secret credentials', () => { - withSourceCacheEnv( - { - [SOURCE_CACHE_CONFIG_ENV]: JSON.stringify({ - endpoint: 'https://cache.example.test', - region: 'us-east-1', - bucket: 'agentconnect-cache', - prefix: '/install-a/', - forcePathStyle: false, - credentialSource: 'secret', - limits: { - maxBundleBytes: 3_000, - orgTotalBytes: 30_000, - pendingReservationMs: 7_200_000, - pendingObjectTtlDays: 3, - unreferencedObjectTtlDays: 8, - unreadPointerTtlDays: 31 - } - }), - [SOURCE_CACHE_ACCESS_KEY_ID_ENV]: 'AKIDEXAMPLE', - [SOURCE_CACHE_SECRET_ACCESS_KEY_ENV]: 'secret-example', - [SOURCE_CACHE_SESSION_TOKEN_ENV]: 'session-example' - }, - () => { - const cfg = loadConfig({ root: tmpRoot() }) - expect(cfg.sourceCache).toEqual({ + it('keeps the feature out of the shared daemon Config and ignores source-cache environment', () => { + expect( + ConfigSchema.parse({ + version: 1, + sourceCache: { endpoint: 'https://cache.example.test', - region: 'us-east-1', bucket: 'agentconnect-cache', - prefix: 'install-a', - forcePathStyle: false, - credentials: { - source: 'secret', - accessKeyId: 'AKIDEXAMPLE', - secretAccessKey: 'secret-example', - sessionToken: 'session-example' - }, - limits: { - maxBundleBytes: 3_000, - orgTotalBytes: 30_000, - pendingReservationMs: 7_200_000, - pendingObjectTtlDays: 3, - unreferencedObjectTtlDays: 8, - unreadPointerTtlDays: 31 + credentials: { source: 'serviceAccount' } + } + }) + ).not.toHaveProperty('sourceCache') + + vi.stubEnv( + 'AC_SOURCE_CACHE_CONFIG', + JSON.stringify({ + endpoint: 'not a url', + bucket: 'agentconnect-cache', + credentialSource: 'serviceAccount' + }) + ) + expect(loadConfig({ root: tmpRoot() })).not.toHaveProperty('sourceCache') + }) + + it('ignores a sourceCache arm in a daemon config file instead of making it a second configuration path', () => { + expect( + loadConfig({ + root: tmpRoot({ + sourceCache: { + endpoint: 'https://cache.example.test', + bucket: 'agentconnect-cache', + credentials: { source: 'serviceAccount' } } }) - } - ) + }) + ).not.toHaveProperty('sourceCache') }) - it("applies the design's defaults, including the ServiceAccount credential source", () => { - withSourceCacheEnv( - { - [SOURCE_CACHE_CONFIG_ENV]: JSON.stringify({ + it('resolves the member document defaults for a ServiceAccount identity', () => { + const root = tmpRoot() + const config = readSourceCacheConfig( + writeDocument( + { + version: 1, endpoint: 'https://cache.example.test', bucket: 'agentconnect-cache', - credentialSource: 'serviceAccount' - }) - }, - () => { - const cfg = loadConfig({ root: tmpRoot() }) - expect(cfg.sourceCache?.credentials).toEqual({ source: 'serviceAccount' }) - expect(cfg.sourceCache?.region).toBe('auto') - expect(cfg.sourceCache?.prefix).toBe('') - expect(cfg.sourceCache?.forcePathStyle).toBe(false) - expect(cfg.sourceCache?.limits).toEqual({ - maxBundleBytes: 2 * 1024 ** 3, - orgTotalBytes: 20 * 1024 ** 3, - pendingReservationMs: 60 * 60 * 1000, - pendingObjectTtlDays: 2, - unreferencedObjectTtlDays: 7, - unreadPointerTtlDays: 30 - }) - } + credentials: { source: 'serviceAccount' } + }, + join(root, 'source-cache.json') + ) ) - }) - it('rejects malformed config rather than disabling silently', () => { - withSourceCacheEnv( - { - [SOURCE_CACHE_CONFIG_ENV]: JSON.stringify({ - endpoint: 'not a url', - bucket: 'agentconnect-cache', - credentialSource: 'serviceAccount' - }) - }, - () => { - expect(() => loadConfig({ root: tmpRoot() })).toThrow(SOURCE_CACHE_CONFIG_ENV) + expect(config).toEqual({ + endpoint: 'https://cache.example.test', + region: 'us-east-1', + bucket: 'agentconnect-cache', + prefix: '', + forcePathStyle: false, + credentials: { source: 'serviceAccount' }, + limits: { + maxBundleBytes: 2 * 1024 ** 3, + orgTotalBytes: 20 * 1024 ** 3, + pendingReservationMs: 60 * 60 * 1000, + pendingObjectTtlDays: 2, + unreferencedObjectTtlDays: 7, + unreadPointerTtlDays: 30 } - ) + }) }) - it('requires both static keys when the credential source is a Secret', () => { - withSourceCacheEnv( - { - [SOURCE_CACHE_CONFIG_ENV]: JSON.stringify({ + it('resolves a two-key Secret document without requiring a session token', () => { + const root = tmpRoot() + const files = writeSecretFiles(root) + const config = readSourceCacheConfig( + writeDocument( + { + version: 1, endpoint: 'https://cache.example.test', + region: 'auto', bucket: 'agentconnect-cache', - credentialSource: 'secret' - }), - [SOURCE_CACHE_ACCESS_KEY_ID_ENV]: 'AKIDEXAMPLE' + prefix: '/install-a/', + forcePathStyle: true, + credentials: { source: 'secret', ...files }, + limits: { maxBundleBytes: 3_000, orgTotalBytes: 30_000 } + }, + join(root, 'source-cache.json') + ) + ) + + expect(config).toEqual({ + endpoint: 'https://cache.example.test', + region: 'auto', + bucket: 'agentconnect-cache', + prefix: 'install-a', + forcePathStyle: true, + credentials: { + source: 'secret', + accessKeyId: 'AKIDEXAMPLE', + secretAccessKey: 'secret-example' }, - () => { - expect(() => loadConfig({ root: tmpRoot() })).toThrow(SOURCE_CACHE_SECRET_ACCESS_KEY_ENV) + limits: { + maxBundleBytes: 3_000, + orgTotalBytes: 30_000, + pendingReservationMs: 60 * 60 * 1000, + pendingObjectTtlDays: 2, + unreferencedObjectTtlDays: 7, + unreadPointerTtlDays: 30 } - ) + }) }) - it('accepts a file document and keeps the environment out of the direct schema', () => { - expect(ConfigSchema.parse({ version: 1 })).not.toHaveProperty('sourceCache') - expect( - ConfigSchema.parse({ - version: 1, - sourceCache: { + it('reads an optional session token only when the document names its file', () => { + const root = tmpRoot() + const files = writeSecretFiles(root, { sessionToken: 'session-example' }) + const config = readSourceCacheConfig( + writeDocument( + { + version: 1, endpoint: 'https://cache.example.test', - region: 'auto', bucket: 'agentconnect-cache', - prefix: '', - forcePathStyle: true, - credentials: { source: 'serviceAccount' }, - limits: {} - } - }).sourceCache - ).toMatchObject({ + credentials: { source: 'secret', ...files } + }, + join(root, 'source-cache.json') + ) + ) + + expect(config?.credentials).toEqual({ + source: 'secret', + accessKeyId: 'AKIDEXAMPLE', + secretAccessKey: 'secret-example', + sessionToken: 'session-example' + }) + }) + + it('rejects an incomplete or malformed member document rather than disabling silently', () => { + const path = writeDocument({ + version: 1, endpoint: 'https://cache.example.test', - forcePathStyle: true, - credentials: { source: 'serviceAccount' }, - limits: { maxBundleBytes: 2 * 1024 ** 3 } + bucket: 'agentconnect-cache', + credentials: { source: 'secret' } + }) + expect(() => readSourceCacheConfig(path)).toThrow(/invalid source cache configuration/i) + expect(() => readSourceCacheConfig(writeDocument({ endpoint: 'not a url' }))).toThrow( + /invalid source cache configuration/i + ) + }) + + it('does not expose inline static credentials or environment-variable arms in the member schema', () => { + const parsed = SourceCacheDocumentSchema.safeParse({ + version: 1, + endpoint: 'https://cache.example.test', + bucket: 'agentconnect-cache', + credentials: { + source: 'secret', + accessKeyId: 'AKIDEXAMPLE', + secretAccessKey: 'secret-example' + } }) + expect(parsed.success).toBe(false) + + expect( + SourceCacheConfigSchema.safeParse({ + endpoint: 'https://cache.example.test', + region: 'auto', + bucket: 'agentconnect-cache', + credentials: { source: 'serviceAccount' }, + limits: {} + }).success + ).toBe(false) + }) + + it('uses fixed member-config paths matching the chart mounts', () => { + expect(SOURCE_CACHE_CONFIG_PATH).toBe('/var/run/ac-source-cache/config.json') + expect(SOURCE_CACHE_ACCESS_KEY_ID_FILE).toBe('/var/run/ac-source-cache-credentials/access-key-id') + expect(SOURCE_CACHE_SECRET_ACCESS_KEY_FILE).toBe('/var/run/ac-source-cache-credentials/secret-access-key') }) }) diff --git a/packages/daemon/test/source-cache-signer.test.ts b/packages/daemon/test/source-cache-signer.test.ts index 6c7af3a83..58d385c3e 100644 --- a/packages/daemon/test/source-cache-signer.test.ts +++ b/packages/daemon/test/source-cache-signer.test.ts @@ -65,18 +65,21 @@ describe('source cache presigner', () => { expect(url.searchParams.get('X-Amz-Signature')).toMatch(/^[0-9a-f]{64}$/) }) - it('allows a caller to shorten a lifetime but never past the SigV4 seven-day ceiling', async () => { + it('fixes GET and PUT at the design lifetimes and exposes no per-call override', async () => { const signer = createSourceCacheSigner(config(), { now: () => NOW })! - const short = new URL(await signer.presignGet('src/org/anon/repo/refs/hash/full/latest', { ttlSeconds: 60 })) - expect(short.searchParams.get('X-Amz-Expires')).toBe('60') - await expect( - signer.presignPut('src/org/anon/repo/bundles/x.bundle', { + const get = new URL(await signer.presignGet('src/org/anon/repo/refs/hash/full/latest')) + expect(get.searchParams.get('X-Amz-Expires')).toBe('300') + + const put = new URL( + await signer.presignPut('src/org/anon/repo/bundles/x.bundle', { contentLength: 1, - checksumSha256: '47DEQpj8HBSa+/TImW+5JCeuQeRkm5NMpJWZG3hSuFU=', - ttlSeconds: 7 * 24 * 60 * 60 + 1 + checksumSha256: '47DEQpj8HBSa+/TImW+5JCeuQeRkm5NMpJWZG3hSuFU=' }) - ).rejects.toThrow(/seven days|604800/i) + ) + expect(put.searchParams.get('X-Amz-Expires')).toBe('900') + expect(SOURCE_CACHE_GET_TTL_SECONDS).toBe(300) + expect(SOURCE_CACHE_PUT_TTL_SECONDS).toBe(900) }) it('uses the Kubernetes service account credential chain without embedding credentials in config', async () => { diff --git a/scripts/test-chart-render.rb b/scripts/test-chart-render.rb index 2f99c02b6..9389634ed 100644 --- a/scripts/test-chart-render.rb +++ b/scripts/test-chart-render.rb @@ -610,14 +610,18 @@ .dig('spec', 'template', 'spec', 'containers', 0).fetch('env').map { |item| item.fetch('name') } abort('an unset grace must leave both sides on the daemon default') if default_member.include?('AC_K8S_ORPHAN_GRACE_MS') -# Source Cache is a daemon-pool member capability only. The default has no bucket configured, so -# it must render no source-cache environment at all; enabling Secret credentials must project -# values by reference, never inline, and must stay out of reconciler/sandbox objects. +# Source Cache is a daemon-pool member capability only. The member reads a fixed document from +# a ConfigMap; static credentials are Secret files, never environment variables or Config arms. +source_cache_name = 'example-agentconnect-source-cache' +default_documents = documents +abort('source Cache must be off by default') if default_documents.any? do |doc| + doc['kind'] == 'ConfigMap' && doc.dig('metadata', 'name') == source_cache_name +end abort('source Cache must be off by default') if container.fetch('env').any? { |item| item.fetch('name').start_with?('AC_SOURCE_CACHE_') } source_cache_set = [ '--set', 'sourceCache.endpoint=https://cache.example.test', - '--set', 'sourceCache.region=us-east-1', + '--set', 'sourceCache.region=auto', '--set', 'sourceCache.bucket=agentconnect-cache', '--set', 'sourceCache.prefix=install-a', '--set', 'sourceCache.forcePathStyle=true', @@ -628,21 +632,23 @@ source_rendered, source_error, source_status = Open3.capture3(*(command + source_cache_set)) abort("helm template (source cache enabled) failed:\n#{source_error}") unless source_status.success? source_documents = YAML.load_stream(source_rendered).compact -source_deployment = source_documents.find do |doc| - doc['kind'] == 'Deployment' && doc.dig('metadata', 'name') == 'example-agentconnect-daemon-pool' -end || abort('missing daemon-pool Deployment in the source-cache render') -source_container = source_deployment.dig('spec', 'template', 'spec', 'containers').find { |item| item['name'] == 'daemon-pool' } || - abort('missing daemon-pool container in the source-cache render') -source_env = source_container.fetch('env') -config_entry = source_env.find { |item| item['name'] == 'AC_SOURCE_CACHE_CONFIG' } || - abort('daemon pool must receive the source-cache config') -source_config = JSON.parse(config_entry.fetch('value')) +source_config_map = source_documents.find do |doc| + doc['kind'] == 'ConfigMap' && doc.dig('metadata', 'name') == source_cache_name +end || abort('daemon pool must receive the source-cache ConfigMap') +source_config = JSON.parse(source_config_map.dig('data', 'config.json')) +abort('source-cache config must carry version 1') unless source_config['version'] == 1 abort('source-cache config must carry the endpoint') unless source_config['endpoint'] == 'https://cache.example.test' -abort('source-cache config must carry the region') unless source_config['region'] == 'us-east-1' +abort('source-cache config must carry the region') unless source_config['region'] == 'auto' abort('source-cache config must carry the bucket') unless source_config['bucket'] == 'agentconnect-cache' abort('source-cache config must carry the prefix') unless source_config['prefix'] == 'install-a' abort('source-cache config must carry forcePathStyle') unless source_config['forcePathStyle'] == true -abort('source-cache config must name the credential source') unless source_config['credentialSource'] == 'secret' +abort('source-cache config must name the credential source') unless source_config.dig('credentials', 'source') == 'secret' +abort('source-cache Secret credentials must identify member-local key files') unless + source_config['credentials'] == { + 'source' => 'secret', + 'accessKeyIdFile' => '/var/run/ac-source-cache-credentials/access-key-id', + 'secretAccessKeyFile' => '/var/run/ac-source-cache-credentials/secret-access-key' + } abort('source-cache config must carry section 10 limits') unless source_config['limits'] == { 'maxBundleBytes' => 3000, 'orgTotalBytes' => 30000, @@ -652,28 +658,60 @@ 'unreadPointerTtlDays' => 31 } -{ - 'AC_SOURCE_CACHE_ACCESS_KEY_ID' => 'AWS_ACCESS_KEY_ID', - 'AC_SOURCE_CACHE_SECRET_ACCESS_KEY' => 'AWS_SECRET_ACCESS_KEY', - 'AC_SOURCE_CACHE_SESSION_TOKEN' => 'AWS_SESSION_TOKEN' -}.each do |env_name, secret_key| - entry = source_env.find { |item| item['name'] == env_name } || abort("daemon pool must project #{env_name}") - abort("#{env_name} must be a Secret reference, never an inline value") if entry.key?('value') - abort("#{env_name} must read #{secret_key} from sourceCache.existingSecret") unless - entry.dig('valueFrom', 'secretKeyRef') == { - 'name' => 'source-cache-credentials', - 'key' => secret_key - } -end +source_deployment = source_documents.find do |doc| + doc['kind'] == 'Deployment' && doc.dig('metadata', 'name') == 'example-agentconnect-daemon-pool' +end || abort('missing daemon-pool Deployment in the source-cache render') +source_pod = source_deployment.dig('spec', 'template', 'spec') +source_container = source_pod.fetch('containers').find { |item| item['name'] == 'daemon-pool' } || + abort('missing daemon-pool container in the source-cache render') +abort('source-cache config must not travel through daemon environment') if + source_container.fetch('env').any? { |item| item.fetch('name').start_with?('AC_SOURCE_CACHE_') } +source_mounts = source_container.fetch('volumeMounts').to_h { |item| [item.fetch('name'), item] } +abort('daemon pool must mount the source-cache ConfigMap') unless + source_mounts.dig('source-cache', 'mountPath') == '/var/run/ac-source-cache' +abort('daemon pool must mount the static credential Secret') unless + source_mounts.dig('source-cache-credentials', 'mountPath') == '/var/run/ac-source-cache-credentials' + +source_volumes = source_pod.fetch('volumes').to_h { |item| [item.fetch('name'), item] } +abort('source-cache ConfigMap volume must name the rendered document') unless + source_volumes.dig('source-cache', 'configMap') == { + 'defaultMode' => 256, + 'items' => [{ 'key' => 'config.json', 'path' => 'config.json' }], + 'name' => source_cache_name + } +abort('source-cache Secret volume must be member-local') unless + source_volumes.dig('source-cache-credentials', 'secret', 'secretName') == 'source-cache-credentials' +abort('two-key static credentials must not require a session-token key') unless + source_volumes.dig('source-cache-credentials', 'secret', 'items') == [ + { 'key' => 'AWS_ACCESS_KEY_ID', 'path' => 'access-key-id' }, + { 'key' => 'AWS_SECRET_ACCESS_KEY', 'path' => 'secret-access-key' } + ] + +source_config_bytes = source_config_map.dig('data', 'config.json') +abort('source-cache ConfigMap must not inline static credentials') if + source_config_bytes.include?('AKIDEXAMPLE') || source_config_bytes.include?('secret-example') source_reconciler = source_documents.find do |doc| doc['kind'] == 'CronJob' && doc.dig('metadata', 'name') == 'example-agentconnect-daemon-pool-reconciler' end || abort('missing daemon-pool reconciler') -abort('the reconciler must not receive source-cache credentials') if source_reconciler.to_yaml.include?('AC_SOURCE_CACHE') +abort('the reconciler must not receive source-cache configuration') if source_reconciler.to_yaml.include?('source-cache') source_runtime = source_documents.find do |doc| doc['kind'] == 'SandboxTemplate' && doc.dig('metadata', 'name') == 'example-agentconnect-runtime' end || abort('missing runtime SandboxTemplate') -abort('sandbox pods must never receive source-cache credentials') if source_runtime.to_yaml.include?('AC_SOURCE_CACHE') +abort('sandbox pods must never receive source-cache configuration') if source_runtime.to_yaml.include?('source-cache') + +session_token_set = source_cache_set + ['--set', 'sourceCache.sessionTokenKey=AWS_SESSION_TOKEN'] +session_rendered, session_error, session_status = Open3.capture3(*(command + session_token_set)) +abort("helm template (source cache session token) failed:\n#{session_error}") unless session_status.success? +session_documents = YAML.load_stream(session_rendered).compact +session_deployment = session_documents.find do |doc| + doc['kind'] == 'Deployment' && doc.dig('metadata', 'name') == 'example-agentconnect-daemon-pool' +end || abort('missing daemon-pool Deployment in the session-token render') +session_volumes = session_deployment.dig('spec', 'template', 'spec', 'volumes').to_h { |item| [item.fetch('name'), item] } +abort('a configured session-token key must be projected as an optional-pair addition') unless + session_volumes.dig('source-cache-credentials', 'secret', 'items').include?( + { 'key' => 'AWS_SESSION_TOKEN', 'path' => 'session-token' } + ) service_account_set = [ '--set', 'sourceCache.endpoint=https://cache.example.test', @@ -683,15 +721,28 @@ service_rendered, service_error, service_status = Open3.capture3(*(command + service_account_set)) abort("helm template (source cache service account) failed:\n#{service_error}") unless service_status.success? service_documents = YAML.load_stream(service_rendered).compact -service_member = service_documents.find do |doc| +service_config_map = service_documents.find do |doc| + doc['kind'] == 'ConfigMap' && doc.dig('metadata', 'name') == source_cache_name +end || abort('missing service-account source-cache ConfigMap') +service_config = JSON.parse(service_config_map.dig('data', 'config.json')) +abort('ServiceAccount source must be AWS identity resolution') unless + service_config.dig('credentials', 'source') == 'serviceAccount' +abort('ServiceAccount source must default to a concrete AWS region') unless service_config['region'] == 'us-east-1' +service_deployment = service_documents.find do |doc| doc['kind'] == 'Deployment' && doc.dig('metadata', 'name') == 'example-agentconnect-daemon-pool' -end.dig('spec', 'template', 'spec', 'containers').find { |item| item['name'] == 'daemon-pool' } -service_env_names = service_member.fetch('env').map { |item| item.fetch('name') } -abort('ServiceAccount source must not inject static-key environment') if service_env_names.any? { |name| name.start_with?('AC_SOURCE_CACHE_ACCESS_KEY', 'AC_SOURCE_CACHE_SECRET') } +end || abort('missing daemon-pool service-account Deployment') +service_pod = service_deployment.dig('spec', 'template', 'spec') +service_container = service_pod.fetch('containers').find { |item| item['name'] == 'daemon-pool' } +service_mount_names = service_container.fetch('volumeMounts').map { |item| item.fetch('name') } +abort('ServiceAccount source must not mount the static credential Secret') if + service_mount_names.include?('source-cache-credentials') +service_env_names = service_container.fetch('env').map { |item| item.fetch('name') } +abort('ServiceAccount source must not inject static-key environment') if service_env_names.any? { |name| name.start_with?('AC_SOURCE_CACHE_') } [ [['sourceCache.endpoint=https://cache.example.test', 'sourceCache.bucket=agentconnect-cache', 'sourceCache.credentialSource=secret'], 'existingSecret'], - [['sourceCache.endpoint=https://cache.example.test'], 'bucket'] + [['sourceCache.endpoint=https://cache.example.test'], 'bucket'], + [['sourceCache.endpoint=https://cache.example.test', 'sourceCache.bucket=agentconnect-cache', 'sourceCache.credentialSource=serviceAccount', 'sourceCache.region=auto'], 'concrete AWS region'] ].each do |settings, expected| extra = settings.flat_map { |setting| ['--set', setting] } _, refused, refused_status = Open3.capture3(*(command + extra)) @@ -699,5 +750,4 @@ abort("refusal for #{settings.join(', ')} must say what is missing:\n#{refused}") unless refused.include?(expected) end - puts 'chart render contract: ok' From 8c20d754da1f28fecf7e48e78791edf0a1a7e4ce Mon Sep 17 00:00:00 2001 From: Harbor404 <2657212322@qq.com> Date: Fri, 2 Oct 2026 13:02:12 +0800 Subject: [PATCH 3/3] test(source-cache): pin MinIO negative cases --- .../daemon/test/source-cache-minio.it.test.ts | 81 +++++++++++++++++++ 1 file changed, 81 insertions(+) diff --git a/packages/daemon/test/source-cache-minio.it.test.ts b/packages/daemon/test/source-cache-minio.it.test.ts index a6443c362..a61282c35 100644 --- a/packages/daemon/test/source-cache-minio.it.test.ts +++ b/packages/daemon/test/source-cache-minio.it.test.ts @@ -60,4 +60,85 @@ describe.skipIf(!ENABLED)('source cache presigner against MinIO', () => { expect(get.status).toBe(200) expect(Buffer.from(await get.arrayBuffer())).toEqual(body) }) + + it('rejects the same signed URL when the body changes the signed Content-Length', async () => { + const signer = createSourceCacheSigner( + SourceCacheConfigSchema.parse({ + endpoint: ENDPOINT, + region: 'us-east-1', + bucket: BUCKET, + prefix: 'it', + forcePathStyle: true, + credentials: { + source: 'secret', + accessKeyId: ACCESS_KEY_ID, + secretAccessKey: SECRET_ACCESS_KEY + }, + limits: {} + }) + )! + const key = `src/contract/${randomUUID()}.bundle` + const body = Buffer.from(`source-cache-presign-length-${randomUUID()}`) + const checksumSha256 = createHash('sha256').update(body).digest('base64') + const putUrl = await signer.presignPut(key, { + contentLength: body.byteLength, + checksumSha256 + }) + const wrongLengthBody = Buffer.concat([body, Buffer.from('x')]) + + const response = await fetch(putUrl, { + method: 'PUT', + headers: { + 'content-length': String(wrongLengthBody.byteLength), + 'x-amz-checksum-sha256': checksumSha256, + 'x-amz-tagging': 'ac-cache=pending' + }, + body: wrongLengthBody + }) + const responseBody = await response.text() + + expect(response.status, responseBody).toBe(403) + expect(responseBody).toContain('SignatureDoesNotMatch') + }) + + it('rejects same-length bytes whose checksum does not match', async () => { + const signer = createSourceCacheSigner( + SourceCacheConfigSchema.parse({ + endpoint: ENDPOINT, + region: 'us-east-1', + bucket: BUCKET, + prefix: 'it', + forcePathStyle: true, + credentials: { + source: 'secret', + accessKeyId: ACCESS_KEY_ID, + secretAccessKey: SECRET_ACCESS_KEY + }, + limits: {} + }) + )! + const key = `src/contract/${randomUUID()}.bundle` + const body = Buffer.from(`source-cache-presign-checksum-${randomUUID()}`) + const checksumSha256 = createHash('sha256').update(body).digest('base64') + const putUrl = await signer.presignPut(key, { + contentLength: body.byteLength, + checksumSha256 + }) + const alteredBody = Buffer.from(body) + alteredBody[0] = alteredBody[0]! ^ 0xff + + const response = await fetch(putUrl, { + method: 'PUT', + headers: { + 'content-length': String(alteredBody.byteLength), + 'x-amz-checksum-sha256': checksumSha256, + 'x-amz-tagging': 'ac-cache=pending' + }, + body: alteredBody + }) + const responseBody = await response.text() + + expect(response.status, responseBody).toBe(400) + expect(responseBody).toContain('XAmzContentChecksumMismatch') + }) })