mirror of
https://github.com/earendil-works/pi.git
synced 2026-06-18 15:54:04 +08:00
fix(coding-agent): narrow pi.dev sync API surface
This commit is contained in:
@@ -1,14 +1,11 @@
|
||||
import { randomUUID } from "node:crypto";
|
||||
import { join } from "node:path";
|
||||
import { AuthStorage } from "../auth-storage.ts";
|
||||
import { getPiDevAuth, PI_DEV_ACTIVITY_SYNC_SCOPE } from "../pi-dev/index.ts";
|
||||
import { PI_DEV_ACTIVITY_SYNC_SCOPE } from "../pi-dev/config.ts";
|
||||
import { PiDevApiError, type PiDevFetch } from "../pi-dev/http.ts";
|
||||
import { getPiDevAuth } from "../pi-dev/oauth.ts";
|
||||
import { SettingsManager } from "../settings-manager.ts";
|
||||
import {
|
||||
ActivitySyncApiError,
|
||||
type ActivitySyncFetch,
|
||||
getActivitySyncWatermark,
|
||||
uploadSessionAnalytics,
|
||||
} from "./api.ts";
|
||||
import { getActivitySyncWatermark, uploadSessionAnalytics } from "./api.ts";
|
||||
import { type ActivitySyncPayload, buildActivitySyncPayloads } from "./payload.ts";
|
||||
import { buildSessionAnalyticsUpload } from "./session-analytics-reader.ts";
|
||||
import {
|
||||
@@ -39,7 +36,7 @@ export interface SyncSessionAnalyticsOptions {
|
||||
sessionsRoot?: string;
|
||||
settingsManager?: SettingsManager;
|
||||
authStorage?: AuthStorage;
|
||||
fetch?: ActivitySyncFetch;
|
||||
fetch?: PiDevFetch;
|
||||
signal?: AbortSignal;
|
||||
now?: Date;
|
||||
}
|
||||
@@ -66,6 +63,22 @@ async function getActivitySyncAccessToken(
|
||||
return auth.available ? auth.accessToken : undefined;
|
||||
}
|
||||
|
||||
async function runWithRefreshRetry<T>(
|
||||
authStorage: AuthStorage,
|
||||
accessToken: string,
|
||||
options: SyncSessionAnalyticsOptions,
|
||||
request: (accessToken: string) => Promise<T>,
|
||||
): Promise<{ value: T; accessToken: string }> {
|
||||
try {
|
||||
return { value: await request(accessToken), accessToken };
|
||||
} catch (error) {
|
||||
if (!(error instanceof PiDevApiError) || error.status !== 401) throw error;
|
||||
const refreshedAccessToken = await getActivitySyncAccessToken(authStorage, options, true);
|
||||
if (!refreshedAccessToken) throw error;
|
||||
return { value: await request(refreshedAccessToken), accessToken: refreshedAccessToken };
|
||||
}
|
||||
}
|
||||
|
||||
async function uploadWithRefreshRetry(
|
||||
authStorage: AuthStorage,
|
||||
accessToken: string,
|
||||
@@ -73,32 +86,18 @@ async function uploadWithRefreshRetry(
|
||||
metadata: { deviceId: string; idempotencyKey: string },
|
||||
options: SyncSessionAnalyticsOptions,
|
||||
): Promise<{ watermark: string; accessToken: string }> {
|
||||
try {
|
||||
const response = await uploadSessionAnalytics({
|
||||
const uploaded = await runWithRefreshRetry(authStorage, accessToken, options, (token) =>
|
||||
uploadSessionAnalytics({
|
||||
fetch: options.fetch,
|
||||
accessToken,
|
||||
accessToken: token,
|
||||
deviceId: metadata.deviceId,
|
||||
watermark: payload.watermark,
|
||||
idempotencyKey: metadata.idempotencyKey,
|
||||
body: payload.body,
|
||||
contentEncoding: payload.contentEncoding,
|
||||
});
|
||||
return { watermark: response.watermark, accessToken };
|
||||
} catch (error) {
|
||||
if (!(error instanceof ActivitySyncApiError) || error.status !== 401) throw error;
|
||||
const refreshedAccessToken = await getActivitySyncAccessToken(authStorage, options, true);
|
||||
if (!refreshedAccessToken) throw error;
|
||||
const response = await uploadSessionAnalytics({
|
||||
fetch: options.fetch,
|
||||
accessToken: refreshedAccessToken,
|
||||
deviceId: metadata.deviceId,
|
||||
watermark: payload.watermark,
|
||||
idempotencyKey: metadata.idempotencyKey,
|
||||
body: payload.body,
|
||||
contentEncoding: payload.contentEncoding,
|
||||
});
|
||||
return { watermark: response.watermark, accessToken: refreshedAccessToken };
|
||||
}
|
||||
}),
|
||||
);
|
||||
return { watermark: uploaded.value.watermark, accessToken: uploaded.accessToken };
|
||||
}
|
||||
|
||||
async function getWatermarkWithRefreshRetry(
|
||||
@@ -107,20 +106,10 @@ async function getWatermarkWithRefreshRetry(
|
||||
deviceId: string,
|
||||
options: SyncSessionAnalyticsOptions,
|
||||
): Promise<{ watermark: string | null; accessToken: string }> {
|
||||
try {
|
||||
const response = await getActivitySyncWatermark(accessToken, deviceId, {
|
||||
fetch: options.fetch,
|
||||
});
|
||||
return { watermark: response.watermark, accessToken };
|
||||
} catch (error) {
|
||||
if (!(error instanceof ActivitySyncApiError) || error.status !== 401) throw error;
|
||||
const refreshedAccessToken = await getActivitySyncAccessToken(authStorage, options, true);
|
||||
if (!refreshedAccessToken) throw error;
|
||||
const response = await getActivitySyncWatermark(refreshedAccessToken, deviceId, {
|
||||
fetch: options.fetch,
|
||||
});
|
||||
return { watermark: response.watermark, accessToken: refreshedAccessToken };
|
||||
}
|
||||
const response = await runWithRefreshRetry(authStorage, accessToken, options, (token) =>
|
||||
getActivitySyncWatermark(token, deviceId, { fetch: options.fetch }),
|
||||
);
|
||||
return { watermark: response.value.watermark, accessToken: response.accessToken };
|
||||
}
|
||||
|
||||
async function syncSessionAnalyticsUnlocked(options: SyncSessionAnalyticsOptions): Promise<ActivitySyncResult> {
|
||||
@@ -148,7 +137,6 @@ async function syncSessionAnalyticsUnlocked(options: SyncSessionAnalyticsOptions
|
||||
});
|
||||
|
||||
if (upload.records.length === 0) {
|
||||
await saveActivitySyncState(state, options.agentDir);
|
||||
return {
|
||||
status: "no_changes",
|
||||
filesScanned: upload.filesScanned,
|
||||
@@ -180,10 +168,11 @@ async function syncSessionAnalyticsUnlocked(options: SyncSessionAnalyticsOptions
|
||||
recordsSent += payload.recordCount;
|
||||
compressedBytes += payload.compressedBytes;
|
||||
decompressedBytes += payload.decompressedBytes;
|
||||
state.lastSuccessAt = (options.now ?? new Date()).toISOString();
|
||||
await saveActivitySyncState(state, options.agentDir);
|
||||
}
|
||||
|
||||
state.lastSuccessAt = (options.now ?? new Date()).toISOString();
|
||||
await saveActivitySyncState(state, options.agentDir);
|
||||
|
||||
return {
|
||||
status: "uploaded",
|
||||
recordsSent,
|
||||
|
||||
@@ -1,37 +1,14 @@
|
||||
import type { Buffer } from "node:buffer";
|
||||
import {
|
||||
PI_DEV_ACTIVITY_SYNC_SCOPE,
|
||||
PI_DEV_DEFAULT_BASE_URL,
|
||||
PI_DEV_OAUTH_CLIENT_ID,
|
||||
PI_DEV_OFFLINE_ACCESS_SCOPE,
|
||||
} from "../pi-dev/config.ts";
|
||||
import {
|
||||
getPiDevApiUrl,
|
||||
getPiDevFetch,
|
||||
isRecord,
|
||||
PiDevApiError,
|
||||
type PiDevApiOptions,
|
||||
type PiDevFetch,
|
||||
readJson,
|
||||
requireNumber,
|
||||
requireString,
|
||||
throwIfPiDevNotOk,
|
||||
} from "../pi-dev/http.ts";
|
||||
import {
|
||||
type PiDevDeviceFlowResponse,
|
||||
type PiDevTokenResponse,
|
||||
pollPiDevDeviceToken,
|
||||
refreshPiDevAccessToken,
|
||||
startPiDevDeviceFlow,
|
||||
} from "../pi-dev/oauth.ts";
|
||||
|
||||
export const ACTIVITY_SYNC_CLIENT_ID = PI_DEV_OAUTH_CLIENT_ID;
|
||||
export const ACTIVITY_SYNC_SCOPE = `${PI_DEV_ACTIVITY_SYNC_SCOPE} ${PI_DEV_OFFLINE_ACCESS_SCOPE}`;
|
||||
export const DEFAULT_PI_DEV_URL = PI_DEV_DEFAULT_BASE_URL;
|
||||
|
||||
export type ActivitySyncDeviceFlowResponse = PiDevDeviceFlowResponse;
|
||||
export type ActivitySyncTokenResponse = PiDevTokenResponse;
|
||||
|
||||
export interface ActivitySyncWatermarkResponse {
|
||||
ok: true;
|
||||
watermark: string | null;
|
||||
@@ -44,9 +21,7 @@ export interface ActivitySyncUploadResponse {
|
||||
watermark: string;
|
||||
}
|
||||
|
||||
export type ActivitySyncFetch = PiDevFetch;
|
||||
|
||||
export interface ActivitySyncApiOptions extends PiDevApiOptions {}
|
||||
export type ActivitySyncApiOptions = PiDevApiOptions;
|
||||
|
||||
export interface UploadSessionAnalyticsOptions extends ActivitySyncApiOptions {
|
||||
accessToken: string;
|
||||
@@ -57,13 +32,6 @@ export interface UploadSessionAnalyticsOptions extends ActivitySyncApiOptions {
|
||||
contentEncoding: "zstd";
|
||||
}
|
||||
|
||||
export class ActivitySyncApiError extends PiDevApiError {
|
||||
constructor(status: number, errorCode: string | undefined, description: string | undefined, operation?: string) {
|
||||
super(status, errorCode, description, operation);
|
||||
this.name = "ActivitySyncApiError";
|
||||
}
|
||||
}
|
||||
|
||||
function parseWatermarkResponse(json: unknown): ActivitySyncWatermarkResponse {
|
||||
if (!isRecord(json) || json.ok !== true || (json.watermark !== null && typeof json.watermark !== "string")) {
|
||||
throw new Error("Invalid activity sync watermark response");
|
||||
@@ -83,38 +51,6 @@ function parseUploadResponse(json: unknown): ActivitySyncUploadResponse {
|
||||
};
|
||||
}
|
||||
|
||||
export async function startActivitySyncDeviceFlow(
|
||||
deviceId: string,
|
||||
options: ActivitySyncApiOptions = {},
|
||||
): Promise<ActivitySyncDeviceFlowResponse> {
|
||||
return startPiDevDeviceFlow({
|
||||
...options,
|
||||
deviceId,
|
||||
scopes: [PI_DEV_ACTIVITY_SYNC_SCOPE],
|
||||
errorClass: ActivitySyncApiError,
|
||||
});
|
||||
}
|
||||
|
||||
export async function pollActivitySyncDeviceToken(
|
||||
deviceCode: string,
|
||||
options: ActivitySyncApiOptions = {},
|
||||
): Promise<ActivitySyncTokenResponse> {
|
||||
return pollPiDevDeviceToken(deviceCode, {
|
||||
...options,
|
||||
errorClass: ActivitySyncApiError,
|
||||
});
|
||||
}
|
||||
|
||||
export async function refreshActivitySyncAccessToken(
|
||||
refreshToken: string,
|
||||
options: ActivitySyncApiOptions = {},
|
||||
): Promise<ActivitySyncTokenResponse> {
|
||||
return refreshPiDevAccessToken(refreshToken, {
|
||||
...options,
|
||||
errorClass: ActivitySyncApiError,
|
||||
});
|
||||
}
|
||||
|
||||
export async function getActivitySyncWatermark(
|
||||
accessToken: string,
|
||||
deviceId: string,
|
||||
@@ -123,7 +59,7 @@ export async function getActivitySyncWatermark(
|
||||
const response = await getPiDevFetch(options.fetch)(getPiDevApiUrl(`/analytics/activity/${deviceId}`), {
|
||||
headers: { Authorization: `Bearer ${accessToken}` },
|
||||
});
|
||||
await throwIfPiDevNotOk(response, "GET /analytics/activity/:deviceId", ActivitySyncApiError);
|
||||
await throwIfPiDevNotOk(response, "GET /analytics/activity/:deviceId");
|
||||
return parseWatermarkResponse(await readJson(response));
|
||||
}
|
||||
|
||||
@@ -141,6 +77,6 @@ export async function uploadSessionAnalytics(
|
||||
},
|
||||
body: options.body,
|
||||
});
|
||||
await throwIfPiDevNotOk(response, "POST /analytics/activity/:deviceId", ActivitySyncApiError);
|
||||
await throwIfPiDevNotOk(response, "POST /analytics/activity/:deviceId");
|
||||
return parseUploadResponse(await readJson(response));
|
||||
}
|
||||
|
||||
@@ -1,5 +1,9 @@
|
||||
/**
|
||||
* Activity sync utilities.
|
||||
* Public activity sync entrypoints used by the agent runtime.
|
||||
*
|
||||
* Keep this barrel narrow. Tests and implementation files should import lower-level
|
||||
* payload/API/state helpers from their leaf modules instead of making them part of
|
||||
* the public activity-sync surface.
|
||||
*/
|
||||
|
||||
export {
|
||||
@@ -8,73 +12,4 @@ export {
|
||||
type SyncSessionAnalyticsOptions,
|
||||
syncSessionAnalytics,
|
||||
} from "./activity-sync.ts";
|
||||
export {
|
||||
ACTIVITY_SYNC_CLIENT_ID,
|
||||
ACTIVITY_SYNC_SCOPE,
|
||||
ActivitySyncApiError,
|
||||
type ActivitySyncApiOptions,
|
||||
type ActivitySyncDeviceFlowResponse,
|
||||
type ActivitySyncFetch,
|
||||
type ActivitySyncTokenResponse,
|
||||
type ActivitySyncUploadResponse,
|
||||
type ActivitySyncWatermarkResponse,
|
||||
DEFAULT_PI_DEV_URL,
|
||||
getActivitySyncWatermark,
|
||||
pollActivitySyncDeviceToken,
|
||||
refreshActivitySyncAccessToken,
|
||||
startActivitySyncDeviceFlow,
|
||||
type UploadSessionAnalyticsOptions,
|
||||
uploadSessionAnalytics,
|
||||
} from "./api.ts";
|
||||
export {
|
||||
ACTIVITY_SYNC_CONTENT_ENCODING,
|
||||
ACTIVITY_SYNC_MAX_COMPRESSED_BYTES,
|
||||
ACTIVITY_SYNC_MAX_DECOMPRESSED_BYTES,
|
||||
type ActivitySyncPayload,
|
||||
type BuildActivitySyncPayloadsOptions,
|
||||
buildActivitySyncPayloads,
|
||||
compareSessionAnalyticsRecords,
|
||||
getSessionAnalyticsRecordTimestamp,
|
||||
serializeSessionAnalyticsNdjson,
|
||||
sortSessionAnalyticsRecords,
|
||||
} from "./payload.ts";
|
||||
export {
|
||||
hashSessionAnalyticsString,
|
||||
type ProjectSessionAnalyticsOptions,
|
||||
type ProjectSessionHeaderAnalyticsOptions,
|
||||
projectSessionEntryForAnalytics,
|
||||
projectSessionForAnalytics,
|
||||
projectSessionHeaderForAnalytics,
|
||||
SESSION_ANALYTICS_SCHEMA_VERSION,
|
||||
type SessionAnalyticsContentStats,
|
||||
type SessionAnalyticsEntryRecord,
|
||||
type SessionAnalyticsRecord,
|
||||
type SessionAnalyticsSessionRecord,
|
||||
type SessionAnalyticsUsage,
|
||||
} from "./session-analytics.ts";
|
||||
export {
|
||||
type BuildSessionAnalyticsUploadOptions,
|
||||
type BuildSessionAnalyticsUploadResult,
|
||||
buildSessionAnalyticsUpload,
|
||||
} from "./session-analytics-reader.ts";
|
||||
export {
|
||||
type DiscoveredSession,
|
||||
type DiscoverSessionFilesOptions,
|
||||
type DiscoverSessionsOptions,
|
||||
discoverSessionFiles,
|
||||
discoverSessions,
|
||||
type SessionDiscoveryPhase,
|
||||
type SessionDiscoveryProgress,
|
||||
type SessionDiscoveryProgressCallback,
|
||||
} from "./session-discovery.ts";
|
||||
export {
|
||||
type ActivitySyncLockResult,
|
||||
type ActivitySyncState,
|
||||
type ActivitySyncStatePaths,
|
||||
getActivitySyncStatePaths,
|
||||
getStableActivitySyncDeviceId,
|
||||
loadActivitySyncState,
|
||||
saveActivitySyncState,
|
||||
updateActivitySyncState,
|
||||
withActivitySyncLock,
|
||||
} from "./state.ts";
|
||||
export { getStableActivitySyncDeviceId, loadActivitySyncState } from "./state.ts";
|
||||
|
||||
@@ -5,79 +5,11 @@
|
||||
export {
|
||||
type ActivitySyncResult,
|
||||
type ActivitySyncStatus,
|
||||
type SyncSessionAnalyticsOptions,
|
||||
syncSessionAnalytics,
|
||||
} from "./activity-sync/activity-sync.ts";
|
||||
export {
|
||||
ACTIVITY_SYNC_CLIENT_ID,
|
||||
ACTIVITY_SYNC_SCOPE,
|
||||
ActivitySyncApiError,
|
||||
type ActivitySyncApiOptions,
|
||||
type ActivitySyncDeviceFlowResponse,
|
||||
type ActivitySyncFetch,
|
||||
type ActivitySyncTokenResponse,
|
||||
type ActivitySyncUploadResponse,
|
||||
type ActivitySyncWatermarkResponse,
|
||||
DEFAULT_PI_DEV_URL,
|
||||
getActivitySyncWatermark,
|
||||
pollActivitySyncDeviceToken,
|
||||
refreshActivitySyncAccessToken,
|
||||
startActivitySyncDeviceFlow,
|
||||
type UploadSessionAnalyticsOptions,
|
||||
uploadSessionAnalytics,
|
||||
} from "./activity-sync/api.ts";
|
||||
export {
|
||||
ACTIVITY_SYNC_CONTENT_ENCODING,
|
||||
ACTIVITY_SYNC_MAX_COMPRESSED_BYTES,
|
||||
ACTIVITY_SYNC_MAX_DECOMPRESSED_BYTES,
|
||||
type ActivitySyncPayload,
|
||||
type BuildActivitySyncPayloadsOptions,
|
||||
buildActivitySyncPayloads,
|
||||
compareSessionAnalyticsRecords,
|
||||
getSessionAnalyticsRecordTimestamp,
|
||||
serializeSessionAnalyticsNdjson,
|
||||
sortSessionAnalyticsRecords,
|
||||
} from "./activity-sync/payload.ts";
|
||||
export {
|
||||
hashSessionAnalyticsString,
|
||||
type ProjectSessionAnalyticsOptions,
|
||||
type ProjectSessionHeaderAnalyticsOptions,
|
||||
projectSessionEntryForAnalytics,
|
||||
projectSessionForAnalytics,
|
||||
projectSessionHeaderForAnalytics,
|
||||
SESSION_ANALYTICS_SCHEMA_VERSION,
|
||||
type SessionAnalyticsContentStats,
|
||||
type SessionAnalyticsEntryRecord,
|
||||
type SessionAnalyticsRecord,
|
||||
type SessionAnalyticsSessionRecord,
|
||||
type SessionAnalyticsUsage,
|
||||
} from "./activity-sync/session-analytics.ts";
|
||||
export {
|
||||
type BuildSessionAnalyticsUploadOptions,
|
||||
type BuildSessionAnalyticsUploadResult,
|
||||
buildSessionAnalyticsUpload,
|
||||
} from "./activity-sync/session-analytics-reader.ts";
|
||||
export {
|
||||
type DiscoveredSession,
|
||||
type DiscoverSessionFilesOptions,
|
||||
type DiscoverSessionsOptions,
|
||||
discoverSessionFiles,
|
||||
discoverSessions,
|
||||
type SessionDiscoveryPhase,
|
||||
type SessionDiscoveryProgress,
|
||||
type SessionDiscoveryProgressCallback,
|
||||
} from "./activity-sync/session-discovery.ts";
|
||||
export {
|
||||
type ActivitySyncLockResult,
|
||||
type ActivitySyncState,
|
||||
type ActivitySyncStatePaths,
|
||||
getActivitySyncStatePaths,
|
||||
getStableActivitySyncDeviceId,
|
||||
loadActivitySyncState,
|
||||
saveActivitySyncState,
|
||||
updateActivitySyncState,
|
||||
withActivitySyncLock,
|
||||
} from "./activity-sync/state.ts";
|
||||
type SyncSessionAnalyticsOptions,
|
||||
syncSessionAnalytics,
|
||||
} from "./activity-sync/index.ts";
|
||||
export {
|
||||
AgentSession,
|
||||
type AgentSessionConfig,
|
||||
@@ -159,71 +91,22 @@ export {
|
||||
type WorkingIndicatorOptions,
|
||||
} from "./extensions/index.ts";
|
||||
export {
|
||||
formatPiDevScopes,
|
||||
getPiDevBaseUrl,
|
||||
normalizePiDevBaseUrl,
|
||||
PI_DEV_ACTIVITY_SYNC_SCOPE,
|
||||
PI_DEV_DEFAULT_BASE_URL,
|
||||
PI_DEV_OAUTH_CLIENT_ID,
|
||||
PI_DEV_OAUTH_PROVIDER_ID,
|
||||
PI_DEV_OFFLINE_ACCESS_SCOPE,
|
||||
formatPiDevShareSuccess,
|
||||
getPiDevAuth,
|
||||
loginPiDev,
|
||||
PI_DEV_PROFILE_CONNECTED_STATUS,
|
||||
PI_DEV_PROFILE_SCOPES,
|
||||
PI_DEV_SESSION_SHARE_SCOPE,
|
||||
PI_DEV_SETUP_PROFILE_CONNECTED_STATUS,
|
||||
scopesFromString,
|
||||
withPiDevOfflineAccess,
|
||||
} from "./pi-dev/config.ts";
|
||||
export {
|
||||
createFormBody,
|
||||
getPiDevApiUrl,
|
||||
getPiDevFetch,
|
||||
isRecord,
|
||||
numberField,
|
||||
PiDevApiError,
|
||||
type PiDevApiErrorCtor,
|
||||
type PiDevApiOptions,
|
||||
type PiDevFetch,
|
||||
readJson,
|
||||
readJsonObject,
|
||||
requireNumber,
|
||||
requireString,
|
||||
stringField,
|
||||
throwIfPiDevNotOk,
|
||||
} from "./pi-dev/http.ts";
|
||||
export {
|
||||
getPiDevAuth,
|
||||
hasPiDevScopes,
|
||||
introspectPiDevAccessToken,
|
||||
loginPiDev,
|
||||
type PiDevAccessIntrospectionResult,
|
||||
type PiDevAuthOptions,
|
||||
type PiDevAuthResult,
|
||||
type PiDevDeviceCodeInfo,
|
||||
type PiDevDeviceFlowOptions,
|
||||
type PiDevDeviceFlowResponse,
|
||||
type PiDevDeviceTokenOptions,
|
||||
type PiDevLoginOptions,
|
||||
type PiDevRefreshTokenOptions,
|
||||
type PiDevTokenResponse,
|
||||
pollPiDevDeviceToken,
|
||||
refreshPiDevAccessToken,
|
||||
startPiDevDeviceFlow,
|
||||
} from "./pi-dev/oauth.ts";
|
||||
export {
|
||||
formatPiDevShareSuccess,
|
||||
formatPiDevShareUploadError,
|
||||
getPiDevShareAuth,
|
||||
loginPiDevShare,
|
||||
type PiDevShareAuthResult,
|
||||
type PiDevShareDeviceAuthInfo,
|
||||
type PiDevShareDeviceAuthOptions,
|
||||
type PiDevShareUploadOptions,
|
||||
type PiDevShareUploadResult,
|
||||
parseShareCommand,
|
||||
type ShareCommandMode,
|
||||
type ShareCommandParseResult,
|
||||
uploadPiDevSessionShare,
|
||||
} from "./pi-dev/session-share.ts";
|
||||
} from "./pi-dev/index.ts";
|
||||
|
||||
export { createSyntheticSourceInfo } from "./source-info.ts";
|
||||
|
||||
@@ -1,63 +1,25 @@
|
||||
/**
|
||||
* Public pi.dev entrypoints used by the agent runtime.
|
||||
*
|
||||
* Keep low-level HTTP/OAuth implementation helpers in their leaf modules so this
|
||||
* barrel does not accidentally turn internal pi.dev plumbing into public API.
|
||||
*/
|
||||
|
||||
export {
|
||||
formatPiDevScopes,
|
||||
getPiDevBaseUrl,
|
||||
normalizePiDevBaseUrl,
|
||||
PI_DEV_ACTIVITY_SYNC_SCOPE,
|
||||
PI_DEV_DEFAULT_BASE_URL,
|
||||
PI_DEV_OAUTH_CLIENT_ID,
|
||||
PI_DEV_OAUTH_PROVIDER_ID,
|
||||
PI_DEV_OFFLINE_ACCESS_SCOPE,
|
||||
PI_DEV_PROFILE_CONNECTED_STATUS,
|
||||
PI_DEV_PROFILE_SCOPES,
|
||||
PI_DEV_SESSION_SHARE_SCOPE,
|
||||
PI_DEV_SETUP_PROFILE_CONNECTED_STATUS,
|
||||
scopesFromString,
|
||||
withPiDevOfflineAccess,
|
||||
} from "./config.ts";
|
||||
export {
|
||||
createFormBody,
|
||||
getPiDevApiUrl,
|
||||
getPiDevFetch,
|
||||
isRecord,
|
||||
numberField,
|
||||
PiDevApiError,
|
||||
type PiDevApiErrorCtor,
|
||||
type PiDevApiOptions,
|
||||
type PiDevFetch,
|
||||
readJson,
|
||||
readJsonObject,
|
||||
requireNumber,
|
||||
requireString,
|
||||
stringField,
|
||||
throwIfPiDevNotOk,
|
||||
} from "./http.ts";
|
||||
export {
|
||||
getPiDevAuth,
|
||||
hasPiDevScopes,
|
||||
introspectPiDevAccessToken,
|
||||
loginPiDev,
|
||||
type PiDevAccessIntrospectionResult,
|
||||
type PiDevAuthOptions,
|
||||
type PiDevAuthResult,
|
||||
type PiDevDeviceCodeInfo,
|
||||
type PiDevDeviceFlowOptions,
|
||||
type PiDevDeviceFlowResponse,
|
||||
type PiDevDeviceTokenOptions,
|
||||
type PiDevLoginOptions,
|
||||
type PiDevRefreshTokenOptions,
|
||||
type PiDevTokenResponse,
|
||||
pollPiDevDeviceToken,
|
||||
refreshPiDevAccessToken,
|
||||
startPiDevDeviceFlow,
|
||||
} from "./oauth.ts";
|
||||
export {
|
||||
formatPiDevShareSuccess,
|
||||
formatPiDevShareUploadError,
|
||||
getPiDevShareAuth,
|
||||
loginPiDevShare,
|
||||
type PiDevShareAuthResult,
|
||||
type PiDevShareDeviceAuthInfo,
|
||||
type PiDevShareDeviceAuthOptions,
|
||||
type PiDevShareUploadOptions,
|
||||
type PiDevShareUploadResult,
|
||||
parseShareCommand,
|
||||
|
||||
@@ -32,7 +32,7 @@ export interface PiDevShareUploadOptions {
|
||||
bytes: Uint8Array;
|
||||
byteSize: number;
|
||||
signal?: AbortSignal;
|
||||
fetchFn?: PiDevFetch;
|
||||
fetch?: PiDevFetch;
|
||||
}
|
||||
|
||||
export function parseShareCommand(text: string): ShareCommandParseResult {
|
||||
@@ -70,7 +70,7 @@ export async function getPiDevShareAuth(authStorage: AuthStorage): Promise<PiDev
|
||||
}
|
||||
|
||||
export async function uploadPiDevSessionShare(options: PiDevShareUploadOptions): Promise<PiDevShareUploadResult> {
|
||||
const response = await getPiDevFetch(options.fetchFn)(getPiDevApiUrl("/api/session-shares"), {
|
||||
const response = await getPiDevFetch(options.fetch)(getPiDevApiUrl("/api/session-shares"), {
|
||||
method: "POST",
|
||||
headers: {
|
||||
Authorization: `Bearer ${options.accessToken}`,
|
||||
|
||||
@@ -10,76 +10,6 @@ export {
|
||||
type SyncSessionAnalyticsOptions,
|
||||
syncSessionAnalytics,
|
||||
} from "./core/activity-sync/activity-sync.ts";
|
||||
export {
|
||||
ACTIVITY_SYNC_CLIENT_ID,
|
||||
ACTIVITY_SYNC_SCOPE,
|
||||
ActivitySyncApiError,
|
||||
type ActivitySyncApiOptions,
|
||||
type ActivitySyncDeviceFlowResponse,
|
||||
type ActivitySyncFetch,
|
||||
type ActivitySyncTokenResponse,
|
||||
type ActivitySyncUploadResponse,
|
||||
type ActivitySyncWatermarkResponse,
|
||||
DEFAULT_PI_DEV_URL,
|
||||
getActivitySyncWatermark,
|
||||
pollActivitySyncDeviceToken,
|
||||
refreshActivitySyncAccessToken,
|
||||
startActivitySyncDeviceFlow,
|
||||
type UploadSessionAnalyticsOptions,
|
||||
uploadSessionAnalytics,
|
||||
} from "./core/activity-sync/api.ts";
|
||||
export {
|
||||
ACTIVITY_SYNC_CONTENT_ENCODING,
|
||||
ACTIVITY_SYNC_MAX_COMPRESSED_BYTES,
|
||||
ACTIVITY_SYNC_MAX_DECOMPRESSED_BYTES,
|
||||
type ActivitySyncPayload,
|
||||
type BuildActivitySyncPayloadsOptions,
|
||||
buildActivitySyncPayloads,
|
||||
compareSessionAnalyticsRecords,
|
||||
getSessionAnalyticsRecordTimestamp,
|
||||
serializeSessionAnalyticsNdjson,
|
||||
sortSessionAnalyticsRecords,
|
||||
} from "./core/activity-sync/payload.ts";
|
||||
export {
|
||||
hashSessionAnalyticsString,
|
||||
type ProjectSessionAnalyticsOptions,
|
||||
type ProjectSessionHeaderAnalyticsOptions,
|
||||
projectSessionEntryForAnalytics,
|
||||
projectSessionForAnalytics,
|
||||
projectSessionHeaderForAnalytics,
|
||||
SESSION_ANALYTICS_SCHEMA_VERSION,
|
||||
type SessionAnalyticsContentStats,
|
||||
type SessionAnalyticsEntryRecord,
|
||||
type SessionAnalyticsRecord,
|
||||
type SessionAnalyticsSessionRecord,
|
||||
type SessionAnalyticsUsage,
|
||||
} from "./core/activity-sync/session-analytics.ts";
|
||||
export {
|
||||
type BuildSessionAnalyticsUploadOptions,
|
||||
type BuildSessionAnalyticsUploadResult,
|
||||
buildSessionAnalyticsUpload,
|
||||
} from "./core/activity-sync/session-analytics-reader.ts";
|
||||
export {
|
||||
type DiscoveredSession,
|
||||
type DiscoverSessionFilesOptions,
|
||||
type DiscoverSessionsOptions,
|
||||
discoverSessionFiles,
|
||||
discoverSessions,
|
||||
type SessionDiscoveryPhase,
|
||||
type SessionDiscoveryProgress,
|
||||
type SessionDiscoveryProgressCallback,
|
||||
} from "./core/activity-sync/session-discovery.ts";
|
||||
export {
|
||||
type ActivitySyncLockResult,
|
||||
type ActivitySyncState,
|
||||
type ActivitySyncStatePaths,
|
||||
getActivitySyncStatePaths,
|
||||
getStableActivitySyncDeviceId,
|
||||
loadActivitySyncState,
|
||||
saveActivitySyncState,
|
||||
updateActivitySyncState,
|
||||
withActivitySyncLock,
|
||||
} from "./core/activity-sync/state.ts";
|
||||
export {
|
||||
AgentSession,
|
||||
type AgentSessionConfig,
|
||||
@@ -246,72 +176,23 @@ export type {
|
||||
} from "./core/package-manager.ts";
|
||||
export { DefaultPackageManager } from "./core/package-manager.ts";
|
||||
export {
|
||||
formatPiDevScopes,
|
||||
getPiDevBaseUrl,
|
||||
normalizePiDevBaseUrl,
|
||||
PI_DEV_ACTIVITY_SYNC_SCOPE,
|
||||
PI_DEV_DEFAULT_BASE_URL,
|
||||
PI_DEV_OAUTH_CLIENT_ID,
|
||||
PI_DEV_OAUTH_PROVIDER_ID,
|
||||
PI_DEV_OFFLINE_ACCESS_SCOPE,
|
||||
formatPiDevShareSuccess,
|
||||
getPiDevAuth,
|
||||
loginPiDev,
|
||||
PI_DEV_PROFILE_CONNECTED_STATUS,
|
||||
PI_DEV_PROFILE_SCOPES,
|
||||
PI_DEV_SESSION_SHARE_SCOPE,
|
||||
PI_DEV_SETUP_PROFILE_CONNECTED_STATUS,
|
||||
scopesFromString,
|
||||
withPiDevOfflineAccess,
|
||||
} from "./core/pi-dev/config.ts";
|
||||
export {
|
||||
createFormBody,
|
||||
getPiDevApiUrl,
|
||||
getPiDevFetch,
|
||||
isRecord,
|
||||
numberField,
|
||||
PiDevApiError,
|
||||
type PiDevApiErrorCtor,
|
||||
type PiDevApiOptions,
|
||||
type PiDevFetch,
|
||||
readJson,
|
||||
readJsonObject,
|
||||
requireNumber,
|
||||
requireString,
|
||||
stringField,
|
||||
throwIfPiDevNotOk,
|
||||
} from "./core/pi-dev/http.ts";
|
||||
export {
|
||||
getPiDevAuth,
|
||||
hasPiDevScopes,
|
||||
introspectPiDevAccessToken,
|
||||
loginPiDev,
|
||||
type PiDevAccessIntrospectionResult,
|
||||
type PiDevAuthOptions,
|
||||
type PiDevAuthResult,
|
||||
type PiDevDeviceCodeInfo,
|
||||
type PiDevDeviceFlowOptions,
|
||||
type PiDevDeviceFlowResponse,
|
||||
type PiDevDeviceTokenOptions,
|
||||
type PiDevLoginOptions,
|
||||
type PiDevRefreshTokenOptions,
|
||||
type PiDevTokenResponse,
|
||||
pollPiDevDeviceToken,
|
||||
refreshPiDevAccessToken,
|
||||
startPiDevDeviceFlow,
|
||||
} from "./core/pi-dev/oauth.ts";
|
||||
export {
|
||||
formatPiDevShareSuccess,
|
||||
formatPiDevShareUploadError,
|
||||
getPiDevShareAuth,
|
||||
loginPiDevShare,
|
||||
type PiDevShareAuthResult,
|
||||
type PiDevShareDeviceAuthInfo,
|
||||
type PiDevShareDeviceAuthOptions,
|
||||
type PiDevShareUploadOptions,
|
||||
type PiDevShareUploadResult,
|
||||
parseShareCommand,
|
||||
type ShareCommandMode,
|
||||
type ShareCommandParseResult,
|
||||
uploadPiDevSessionShare,
|
||||
} from "./core/pi-dev/session-share.ts";
|
||||
} from "./core/pi-dev/index.ts";
|
||||
export type {
|
||||
ResourceCollision,
|
||||
ResourceDiagnostic,
|
||||
|
||||
@@ -1,12 +1,7 @@
|
||||
import { Buffer } from "node:buffer";
|
||||
import { afterEach, describe, expect, it, vi } from "vitest";
|
||||
import {
|
||||
type ActivitySyncApiError,
|
||||
getActivitySyncWatermark,
|
||||
refreshActivitySyncAccessToken,
|
||||
startActivitySyncDeviceFlow,
|
||||
uploadSessionAnalytics,
|
||||
} from "../src/core/activity-sync/api.ts";
|
||||
import { getActivitySyncWatermark, uploadSessionAnalytics } from "../src/core/activity-sync/api.ts";
|
||||
import type { PiDevApiError } from "../src/core/pi-dev/http.ts";
|
||||
|
||||
function jsonResponse(body: unknown, status = 200): Response {
|
||||
return new Response(JSON.stringify(body), { status, headers: { "Content-Type": "application/json" } });
|
||||
@@ -17,66 +12,22 @@ afterEach(() => {
|
||||
});
|
||||
|
||||
describe("activity sync api", () => {
|
||||
it("starts the OAuth device flow with the activity sync scope", async () => {
|
||||
vi.stubEnv("PI_DEV_URL", "https://example.test/");
|
||||
let request: Request | undefined;
|
||||
const fetchMock: typeof fetch = async (input, init) => {
|
||||
request = new Request(input, init);
|
||||
return jsonResponse({
|
||||
device_code: "pigd_1",
|
||||
user_code: "ABCD-EFGH",
|
||||
verification_uri: "https://pi.dev/pair",
|
||||
verification_uri_complete: "https://pi.dev/pair?code=ABCD-EFGH",
|
||||
expires_in: 300,
|
||||
interval: 2,
|
||||
});
|
||||
};
|
||||
|
||||
const response = await startActivitySyncDeviceFlow("00000000-0000-4000-8000-000000000000", {
|
||||
fetch: fetchMock,
|
||||
});
|
||||
|
||||
expect(response.device_code).toBe("pigd_1");
|
||||
expect(request?.url).toBe("https://example.test/api/oauth/device");
|
||||
expect(request?.method).toBe("POST");
|
||||
const body = new URLSearchParams(await request?.text());
|
||||
expect(body.get("client_id")).toBe("pi-coding-agent");
|
||||
expect(body.get("scope")).toBe("activity_sync offline_access");
|
||||
expect(body.get("device_id")).toBe("00000000-0000-4000-8000-000000000000");
|
||||
});
|
||||
|
||||
it("refreshes tokens and reads watermarks", async () => {
|
||||
it("reads watermarks", async () => {
|
||||
vi.stubEnv("PI_DEV_URL", "https://example.test");
|
||||
const urls: string[] = [];
|
||||
const fetchMock: typeof fetch = async (input, init) => {
|
||||
const request = new Request(input, init);
|
||||
urls.push(request.url);
|
||||
if (request.url.endsWith("/api/oauth/token")) {
|
||||
return jsonResponse({
|
||||
token_type: "Bearer",
|
||||
access_token: "access-2",
|
||||
refresh_token: "refresh-2",
|
||||
expires_in: 86400,
|
||||
scope: "activity_sync offline_access",
|
||||
});
|
||||
}
|
||||
expect(request.headers.get("Authorization")).toBe("Bearer access-2");
|
||||
return jsonResponse({ ok: true, watermark: "2026-01-01T00:00:00.000Z" });
|
||||
};
|
||||
|
||||
const token = await refreshActivitySyncAccessToken("refresh-1", {
|
||||
fetch: fetchMock,
|
||||
});
|
||||
const watermark = await getActivitySyncWatermark(token.access_token, "device-1", {
|
||||
const watermark = await getActivitySyncWatermark("access-2", "device-1", {
|
||||
fetch: fetchMock,
|
||||
});
|
||||
|
||||
expect(token.refresh_token).toBe("refresh-2");
|
||||
expect(watermark.watermark).toBe("2026-01-01T00:00:00.000Z");
|
||||
expect(urls).toEqual([
|
||||
"https://example.test/api/oauth/token",
|
||||
"https://example.test/analytics/activity/device-1",
|
||||
]);
|
||||
expect(urls).toEqual(["https://example.test/analytics/activity/device-1"]);
|
||||
});
|
||||
|
||||
it("uploads compressed NDJSON with sync headers and surfaces API errors", async () => {
|
||||
@@ -125,12 +76,12 @@ describe("activity sync api", () => {
|
||||
contentEncoding: "zstd",
|
||||
}),
|
||||
).rejects.toMatchObject({
|
||||
name: "ActivitySyncApiError",
|
||||
name: "PiDevApiError",
|
||||
status: 400,
|
||||
errorCode: "invalid_payload",
|
||||
description: "bad line",
|
||||
operation: "POST /analytics/activity/:deviceId",
|
||||
message: "POST /analytics/activity/:deviceId failed: invalid_payload: bad line",
|
||||
} satisfies Partial<ActivitySyncApiError>);
|
||||
} satisfies Partial<PiDevApiError>);
|
||||
});
|
||||
});
|
||||
|
||||
@@ -3,16 +3,15 @@ import { tmpdir } from "node:os";
|
||||
import { join } from "node:path";
|
||||
import { afterEach, describe, expect, it, vi } from "vitest";
|
||||
import { AuthStorage } from "../src/core/auth-storage.ts";
|
||||
import { getPiDevBaseUrl, PI_DEV_SESSION_SHARE_SCOPE } from "../src/core/pi-dev/config.ts";
|
||||
import { getPiDevAuth } from "../src/core/pi-dev/oauth.ts";
|
||||
import {
|
||||
formatPiDevShareSuccess,
|
||||
getPiDevAuth,
|
||||
getPiDevBaseUrl,
|
||||
getPiDevShareAuth,
|
||||
loginPiDevShare,
|
||||
PI_DEV_SESSION_SHARE_SCOPE,
|
||||
parseShareCommand,
|
||||
uploadPiDevSessionShare,
|
||||
} from "../src/core/pi-dev/index.ts";
|
||||
} from "../src/core/pi-dev/session-share.ts";
|
||||
|
||||
const tempDirs: string[] = [];
|
||||
|
||||
@@ -197,7 +196,7 @@ describe("session share client", () => {
|
||||
accessToken: "piga_share",
|
||||
bytes,
|
||||
byteSize: bytes.byteLength,
|
||||
fetchFn,
|
||||
fetch: fetchFn,
|
||||
});
|
||||
|
||||
expect(result).toEqual({ id: "psh_123", url: "https://pi.dev/session/#pi/psh_123" });
|
||||
|
||||
Reference in New Issue
Block a user