@agentium/harness. Types such as ChatMessage, RunContext, and ExecutionServices come from @agentium/core. A ? marks an optional field.
defineWatch
/** Pure definition: validates/copies JSON, without binding services or activating anything. */
export declare function defineWatch(input: WatchDefinitionInput): WatchDefinition;
describeWatch
export declare function describeWatch(definition: WatchDefinitionInput): {
definition: WatchDefinition;
hash: string;
activation: "explicit";
modelCalls: false;
};
DurableWatch
/** A bounded, versioned watch. Construction is inert; all I/O requires an explicit method call. */
export declare class DurableWatch {
readonly definition: WatchDefinition;
readonly key: {
tenantId: string;
taskId: string;
};
constructor(definition: WatchDefinitionInput, services: WatchServices);
activate(): Promise<void>;
resume(): Promise<void>;
inspect(): Promise<WatchState>;
pause(): Promise<void>;
delete(): Promise<void>;
/** Explicit fail-closed rollout. Old pending/unknown sends stay bound to the old definition. */
update(definition: WatchDefinitionInput, services?: WatchServices): Promise<DurableWatch>;
trigger(raw: unknown): Promise<void>;
poll(): Promise<void>;
renew(): Promise<void>;
/** Delivers at most one digest. It never retries an unresolved external effect. */
flush(): Promise<void>;
/** Reconciliation remains available after pause/delete/update; it never dispatches a send. */
reconcile(outboxId: string): Promise<void>;
}
WatchIdentity
export interface WatchIdentity {
tenantId: string;
actorId: string;
}
WatchQuietHours
export interface WatchQuietHours {
start: string;
end: string;
}
WatchLimits
export interface WatchLimits {
maxOperations: number;
maxWakes: number;
maxEventsPerWake: number;
maxPagesPerWake: number;
maxResyncEvents: number;
maxPendingEvents: number;
maxDecisions: number;
maxOutbox: number;
maxEventBytes: number;
maxStateBytes: number;
maxNotificationsPerDay: number;
}
WatchDefinition
export interface WatchDefinition {
id: string;
version: number;
identity: WatchIdentity;
sourceId: string;
sourceScope: string;
destination: string;
channel: string;
policyRevision: number;
grantRefs: string[];
timeZone: string;
quietHours?: WatchQuietHours;
pollIntervalMs: number;
renewalIntervalMs: number;
cooldownMs: number;
limits: WatchLimits;
}
WatchDefinitionInput
export type WatchDefinitionInput = Omit<WatchDefinition, "limits" | "pollIntervalMs" | "renewalIntervalMs" | "cooldownMs"> & {
limits?: Partial<WatchLimits>;
pollIntervalMs?: number;
renewalIntervalMs?: number;
cooldownMs?: number;
};
WatchEvent
export interface WatchEvent {
id: string;
occurredAt: number;
data: DurableJSON;
}
WatchSubscription
export interface WatchSubscription {
reference: string;
expiresAt: number;
}
WatchSourceBatch
export interface WatchSourceBatch {
cursor: string;
events: WatchEvent[];
resynced: boolean;
}
VerifiedWatchTrigger
export interface VerifiedWatchTrigger extends WatchIdentity {
sourceId: string;
sourceScope: string;
eventId: string;
cursor: string;
}
WatchSource
/** A scoped, host-owned source. Implementations must honor read bounds before returning a cursor. */
export interface WatchSource {
id: string;
scope: string;
capabilities: {
idempotentActivation: true;
polling: true;
push: boolean;
};
baseline(signal: AbortSignal): Promise<string>;
activate(key: string, signal: AbortSignal): Promise<WatchSubscription>;
stop(subscription: WatchSubscription | null, signal: AbortSignal): Promise<void>;
read(cursor: string, limits: Pick<WatchLimits, "maxEventsPerWake" | "maxPagesPerWake" | "maxResyncEvents">, signal: AbortSignal): Promise<WatchSourceBatch>;
compareCursors(left: string, right: string): number;
/** Verify signature/audience and authenticated mailbox ownership before returning a trusted hint. */
verifyTrigger?(raw: unknown): Promise<VerifiedWatchTrigger>;
}
WatchScheduleKind
export type WatchScheduleKind = "poll" | "renew" | "flush";
WatchSchedule
export interface WatchSchedule {
key: string;
identity: WatchIdentity;
watchId: string;
configVersion: number;
kind: WatchScheduleKind;
at: number;
}
WatchScheduler
export interface WatchScheduler {
capabilities: {
durable: boolean;
idempotentUpsert: true;
};
schedule(job: WatchSchedule): Promise<void>;
cancel(key: string): Promise<void>;
}
WatchOperation
export type WatchOperation = "activate" | "update" | "pause" | "resume" | "delete" | "trigger" | "poll" | "renew" | "flush" | "reconcile" | "inspect";
WatchAuthorization
export interface WatchAuthorization {
operation: WatchOperation;
definition: Readonly<WatchDefinition>;
/** Notification count includes digests; authorization must validate destination/channel and amount. */
notifications: number;
events: number;
}
WatchServices
export interface WatchServices {
store: DurableTaskStore;
source: WatchSource;
scheduler: WatchScheduler;
notifications: DurableActionConnector & {
channel: string;
destination: string;
};
/** Trusted host authority, including exclusive ownership of a source mailbox where required. */
authorize(request: WatchAuthorization): Promise<boolean>;
/** Deterministic local classification only; no model/network work in this callback. */
filter?(event: Readonly<WatchEvent>): boolean;
leaseMs?: number;
requireDurability?: boolean;
}
WatchOutbox
export interface WatchOutbox {
id: string;
eventIds: string[];
events: WatchEvent[];
state: "pending" | "confirmed";
reservedDay: string | null;
resultRef: string | null;
}
WatchState
export interface WatchState {
status: "inactive" | "activating" | "active" | "paused" | "deleted";
cursor: string | null;
subscription: WatchSubscription | null;
wakes: number;
decisions: Record<string, {
id: string;
decision: "pending" | "ignored" | "resync" | "outbox";
at: number;
}>;
pending: WatchEvent[];
outbox: Record<string, WatchOutbox>;
daily: Record<string, number>;
cooldownUntil: number;
schedules: Partial<Record<WatchScheduleKind, number>>;
}