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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
10 changes: 5 additions & 5 deletions src/lib/cron.ts
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,7 @@ import { generateProjections } from './projections';
import { discoverRevenueSources } from './revenue';
import { generatePaymentPlan, savePaymentPlan } from './payment-planner';
import { reconcileNotionDisputes } from './dispute-sync';
import { enqueueJob, processQueue } from './job-dispatcher';
import { enqueueJob, processQueue, type ScrapeJobType } from './job-dispatcher';

/**
* Cron sync orchestrator.
Expand Down Expand Up @@ -146,7 +146,7 @@ export async function runCronSync(
await enqueueJob(sql, 'portal_scrape', { portal: target }, {
chittyId,
cronSource: 'utility_scrape',
});
}, env);
} catch (err) {
console.error(`[cron:utility:${target}] enqueue failed:`, err);
}
Expand Down Expand Up @@ -175,7 +175,7 @@ export async function runCronSync(
await enqueueJob(sql, 'court_docket', { case_number: '2024D007847' }, {
chittyId,
cronSource: 'court_docket',
});
}, env);
const queueResult = await processQueue(sql, env, ctx);
recordsSynced += queueResult.succeeded;
console.log(`[cron:court_docket] dispatcher: ${queueResult.succeeded} succeeded, ${queueResult.failed} failed`);
Expand Down Expand Up @@ -581,7 +581,7 @@ async function syncMonthlyChecksViaDispatcher(
await enqueueJob(sql, 'mr_cooper', { property: 'addison' }, {
chittyId,
cronSource: 'monthly_check',
});
}, env);
} catch (err) {
console.error('[cron:mr_cooper] enqueue failed:', err);
}
Expand All @@ -596,7 +596,7 @@ async function syncMonthlyChecksViaDispatcher(
}, {
chittyId,
cronSource: 'monthly_check',
});
}, env);
}
} catch (err) {
console.error('[cron:cook_county_tax] enqueue failed:', err);
Expand Down
196 changes: 0 additions & 196 deletions src/lib/fan-out.ts

This file was deleted.

113 changes: 113 additions & 0 deletions src/lib/integrations.ts
Original file line number Diff line number Diff line change
Expand Up @@ -312,6 +312,32 @@ export function connectClient(env: Env) {
const baseUrl = env.CHITTYCONNECT_URL;
if (!baseUrl) return null;

async function connectPost<T>(path: string, body: unknown): Promise<T | null> {
try {
const headers: Record<string, string> = {
'Content-Type': 'application/json',
'X-Source-Service': 'chittycommand',
};
if (env.CHITTY_CONNECT_TOKEN) {
headers['Authorization'] = `Bearer ${env.CHITTY_CONNECT_TOKEN}`;
}
Comment thread
chitcommit marked this conversation as resolved.
const res = await fetch(`${baseUrl}${path}`, {
method: 'POST',
headers,
body: JSON.stringify(body),
signal: AbortSignal.timeout(30000),
});
if (!res.ok) {
console.error(`[connect] POST ${path} failed: ${res.status}`);
return null;
}
return await res.json() as T;
} catch (err) {
console.error(`[connect] POST ${path} error:`, err);
return null;
}
}

return {
/** Discover a service URL by name */
discover: async (serviceName: string): Promise<string | null> => {
Expand All @@ -338,9 +364,48 @@ export function connectClient(env: Env) {
return data.url;
} catch { return null; }
},

// ── Prompt Registry (ContextConsciousness) ─────────────────
/** Resolve a prompt: compose base + layers, apply env gating */
resolvePrompt: (promptId: string, environment: string, variables?: Record<string, string>, additionalLayers?: string[]) =>
connectPost<PromptResolveResponse>('/api/v1/context/prompts/resolve', {
promptId,
environment,
variables,
additionalLayers,
consumerService: 'chittycommand',
}),

/** Execute a prompt: resolve + dispatch to agent, return AI result */
executePrompt: (promptId: string, environment: string, input: Record<string, unknown>, opts?: { additionalLayers?: string[] }) =>
connectPost<PromptExecuteResponse>('/api/v1/context/prompts/execute', {
promptId,
environment,
input,
additionalLayers: opts?.additionalLayers,
consumerService: 'chittycommand',
}),
};
}

export interface PromptResolveResponse {
systemPrompt: string;
aiEnabled: boolean;
version: number;
resolvedLayers: string[];
fallbackMode: string | null;
}

export interface PromptExecuteResponse {
result: string;
promptVersion: number;
resolvedLayers: string[];
executedBy: string;
latencyMs: number;
executionId: number;
aiEnabled: boolean;
}

// ── Mercury ─────────────────────────────────────────────────
// Direct Mercury API for multi-entity banking

Expand Down Expand Up @@ -627,9 +692,57 @@ export function routerClient(env: Env) {
labels: string[];
reasoning?: string;
}>('/agents/triage/classify', payload),

// ── ScrapeAgent proxy methods ──────────────────────────────
/** Enqueue a scrape job on ChittyRouter ScrapeAgent */
enqueueScrapeJob: (jobType: string, target: Record<string, unknown>, opts?: { chittyId?: string; maxAttempts?: number; cronSource?: string }) =>
post<{ id: string; status: string }>('/agents/scrape/enqueue', { jobType, target, ...opts }),

/** Get a single scrape job status */
getScrapeJobStatus: (jobId: string) =>
get<ScrapeJobResponse>(`/agents/scrape/jobs/${encodeURIComponent(jobId)}`),

/** List scrape jobs with filters */
listScrapeJobs: (filters?: { status?: string; jobType?: string; limit?: number }) => {
const params = new URLSearchParams();
if (filters?.status) params.set('status', filters.status);
if (filters?.jobType) params.set('jobType', filters.jobType);
if (filters?.limit) params.set('limit', String(filters.limit));
const qs = params.toString();
return get<{ jobs: ScrapeJobResponse[]; total: number }>(`/agents/scrape/jobs${qs ? `?${qs}` : ''}`);
},

/** Retry a failed/dead-lettered scrape job */
retryScrapeJob: (jobId: string) =>
post<{ status: string }>(`/agents/scrape/jobs/${encodeURIComponent(jobId)}/retry`, {}),

/** Get dead-lettered scrape jobs */
getScrapeDeadLetters: () =>
get<{ jobs: ScrapeJobResponse[] }>('/agents/scrape/dead-letters'),

/** Trigger queue processing on ScrapeAgent */
processScrapeQueue: () =>
post<{ processed: number; succeeded: number; failed: number }>('/agents/scrape/process', {}),

/** Get ScrapeAgent health status */
getScrapeStatus: () =>
get<Record<string, unknown>>('/agents/scrape/status'),
};
}

export interface ScrapeJobResponse {
id: string;
jobType: string;
target: Record<string, unknown>;
status: string;
attempt: number;
maxAttempts: number;
result?: Record<string, unknown>;
error?: string;
createdAt: string;
completedAt?: string;
}

// ── Notion (write path) ───────────────────────────────────────
// Reading Notion is handled by syncNotionTasks() in cron.ts.
// This client covers the write path: creating task pages from disputes.
Expand Down
Loading
Loading