({ host, jsonEmitter, setStreamRequestId }: StdinStreamModeOptions)
| 340 | } |
| 341 | |
| 342 | export async function runStdinStreamMode({ host, jsonEmitter, setStreamRequestId }: StdinStreamModeOptions) { |
| 343 | let hasReceivedStdinCommand = false |
| 344 | let shouldShutdown = false |
| 345 | let activeTaskPromise: Promise<void> | null = null |
| 346 | let fatalStreamError: Error | null = null |
| 347 | let activeRequestId: string | undefined |
| 348 | let activeTaskCommand: "start" | undefined |
| 349 | let latestTaskId: string | undefined |
| 350 | let cancelRequestedForActiveTask = false |
| 351 | let awaitingPostCancelRecovery = false |
| 352 | let hasSeenQueueState = false |
| 353 | let lastQueueDepth = 0 |
| 354 | let lastQueueMessageIds: string[] = [] |
| 355 | const pendingQueuedMessageRequestIds: string[] = [] |
| 356 | const queueMessageRequestIdByMessageId = new Map<string, string>() |
| 357 | |
| 358 | const assignRequestIdsToNewQueueMessages = (queueMessageIds: string[]) => { |
| 359 | for (const messageId of queueMessageIds) { |
| 360 | if (queueMessageRequestIdByMessageId.has(messageId)) { |
| 361 | continue |
| 362 | } |
| 363 | |
| 364 | const requestId = pendingQueuedMessageRequestIds.shift() |
| 365 | if (!requestId) { |
| 366 | continue |
| 367 | } |
| 368 | |
| 369 | queueMessageRequestIdByMessageId.set(messageId, requestId) |
| 370 | } |
| 371 | } |
| 372 | |
| 373 | const promoteRequestIdForDequeuedMessages = (queueMessageIds: string[]) => { |
| 374 | if (lastQueueMessageIds.length === 0) { |
| 375 | return |
| 376 | } |
| 377 | |
| 378 | const remainingIds = new Set(queueMessageIds) |
| 379 | |
| 380 | for (const dequeuedMessageId of lastQueueMessageIds) { |
| 381 | if (remainingIds.has(dequeuedMessageId)) { |
| 382 | continue |
| 383 | } |
| 384 | |
| 385 | const requestId = queueMessageRequestIdByMessageId.get(dequeuedMessageId) |
| 386 | if (requestId) { |
| 387 | setStreamRequestId(requestId) |
| 388 | } |
| 389 | queueMessageRequestIdByMessageId.delete(dequeuedMessageId) |
| 390 | } |
| 391 | } |
| 392 | |
| 393 | const waitForPreviousTaskToSettle = async () => { |
| 394 | if (!activeTaskPromise) { |
| 395 | return |
| 396 | } |
| 397 | |
| 398 | try { |
| 399 | await activeTaskPromise |
no test coverage detected