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
49 changes: 34 additions & 15 deletions src/codex/middle/ConversationProjection.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -272,20 +272,32 @@ struct ProjectedNode {
VisibleCardData card;
};

std::optional<std::size_t> admissionBoundaryPosition(
const std::optional<AuthoritativeItemKey> &admissionAnchor,
bool admissionAtStart, const AuthoritativeItemIndex &authoritativeItems) {
if (admissionAnchor) {
const auto anchor = authoritativeItems.position(*admissionAnchor);
if (anchor)
return (*anchor + 1) * 2;
}
if (admissionAtStart)
return 0;
return std::nullopt;
}

std::size_t submissionPosition(
const PromptSubmission &submission,
const AuthoritativeItemIndex &authoritativeItems,
std::optional<std::size_t> materializedIndex = std::nullopt) {
if (submission.admissionAnchor) {
const auto anchor =
authoritativeItems.position(*submission.admissionAnchor);
if (anchor)
return (*anchor + 1) * 2;
}
const auto admitted = admissionBoundaryPosition(submission.admissionAnchor,
submission.admissionAtStart,
authoritativeItems);
if (admitted)
return *admitted;
if (materializedIndex)
return *materializedIndex * 2 + 1;
// No authoritative tail was known at admission. Until reconcile establishes
// one, the prompt is a tail item rather than a synthetic history prefix.
// A queued pre-hydration prompt has no committed boundary yet. Until
// reconcile establishes one, keep it at the tail of retained history.
return authoritativeItems.ordered.size() * 2 + 2;
}

Expand Down Expand Up @@ -322,15 +334,22 @@ ConversationSnapshot ConversationProjection::project(
binding->second->localCardVisible(nowMilliseconds))
continue;
CardKey visualKey =
item.localPromptKey ? CardKey{*item.localPromptKey} : CardKey{item.key};
item.promptAlias ? CardKey{item.promptAlias->key} : CardKey{item.key};
if (binding != bindings.end())
visualKey = LocalPromptKey{binding->second->id};
const std::size_t position =
binding == bindings.end()
? index * 2 + 1
: submissionPosition(*binding->second, authoritativeItems, index);
const std::uint64_t tieBreaker =
binding == bindings.end() ? 0 : binding->second->admissionOrdinal;
std::size_t position = index * 2 + 1;
std::uint64_t tieBreaker = 0;
if (binding != bindings.end()) {
position =
submissionPosition(*binding->second, authoritativeItems, index);
tieBreaker = binding->second->admissionOrdinal;
} else if (item.promptAlias) {
const auto admitted = admissionBoundaryPosition(
item.promptAlias->admissionAnchor, !item.promptAlias->admissionAnchor,
authoritativeItems);
position = admitted.value_or(position);
tieBreaker = item.promptAlias->admissionOrdinal;
}
VisibleCardData card =
authoritativeCard(item.key, *item.presentation, std::move(visualKey));
nodes.push_back({position, tieBreaker,
Expand Down
23 changes: 16 additions & 7 deletions src/codex/middle/PromptCoordinator.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -172,6 +172,7 @@ PromptCoordinator::beginNext(const std::string &threadId,
});
if (next == found->second.end())
return std::nullopt;
next->admissionAtStart = !next->admissionAnchor;
next->state = PromptState::InFlight;
// Start versus steer is an operation-time fact. A turn which was active
// when the prompt entered the local queue may have completed meanwhile.
Expand Down Expand Up @@ -216,6 +217,7 @@ bool PromptCoordinator::requeue(const std::string &threadId,
if (!pending || pending->state != PromptState::InFlight)
return false;
pending->state = PromptState::Queued;
pending->admissionAtStart = false;
return true;
}

Expand Down Expand Up @@ -246,10 +248,12 @@ bool PromptCoordinator::reassignThread(const std::string &fromThreadId,
auto moved = std::move(aliases->second);
visualAliasesByThread.erase(aliases);
auto &target = visualAliasesByThread[toThreadId];
for (const auto &[key, localKey] : moved) {
for (auto &[key, alias] : moved) {
AuthoritativeItemKey reassigned = key;
reassigned.threadId = toThreadId;
target.insert_or_assign(std::move(reassigned), localKey);
if (alias.admissionAnchor)
alias.admissionAnchor->threadId = toThreadId;
target.insert_or_assign(std::move(reassigned), std::move(alias));
}
};
auto source = byThread.find(fromThreadId);
Expand Down Expand Up @@ -307,7 +311,7 @@ void PromptCoordinator::reconcile(const std::string &threadId,
std::vector<bool> claimed(authoritativeItems.ordered.size());
for (std::size_t index = 0; index < authoritativeItems.ordered.size();
++index)
if (authoritativeItems.ordered[index].localPromptKey)
if (authoritativeItems.ordered[index].promptAlias)
claimed[index] = true;
for (const PromptSubmission &submission : found->second)
if (submission.materializedItem) {
Expand All @@ -322,8 +326,10 @@ void PromptCoordinator::reconcile(const std::string &threadId,
continue;
if (!submission.admissionAnchor &&
submission.state == PromptState::Queued &&
!authoritativeItems.ordered.empty())
!authoritativeItems.ordered.empty()) {
submission.admissionAnchor = authoritativeItems.ordered.back().key;
submission.admissionAtStart = false;
}

const auto exact = authoritativeItems.userMessagesByClientId.find(
submission.clientUserMessageId);
Expand Down Expand Up @@ -381,7 +387,10 @@ void PromptCoordinator::reconcile(const std::string &threadId,
submission.acceptedTransitionActive(nowMilliseconds))
continue;
visualAliasesByThread[threadId].insert_or_assign(
*submission.materializedItem, LocalPromptKey{submission.id});
*submission.materializedItem,
PromptVisualAlias{LocalPromptKey{submission.id},
submission.admissionAnchor,
submission.admissionOrdinal});
}
std::erase_if(found->second, [nowMilliseconds](const auto &submission) {
return submission.state == PromptState::Accepted &&
Expand Down Expand Up @@ -441,10 +450,10 @@ void PromptCoordinator::applyVisualAliases(
const auto aliases = visualAliasesByThread.find(threadId);
if (aliases == visualAliasesByThread.end())
return;
for (const auto &[key, localKey] : aliases->second) {
for (const auto &[key, alias] : aliases->second) {
const auto position = authoritativeItems.position(key);
if (position)
authoritativeItems.ordered[*position].localPromptKey = localKey;
authoritativeItems.ordered[*position].promptAlias = alias;
}
}

Expand Down
13 changes: 10 additions & 3 deletions src/codex/middle/PromptCoordinator.h
Original file line number Diff line number Diff line change
Expand Up @@ -40,6 +40,7 @@ struct PromptSubmission {
qint64 acceptedAtMilliseconds = 0;
QString error;
std::optional<AuthoritativeItemKey> admissionAnchor;
bool admissionAtStart = false;
std::optional<std::string> expectedTurnId;
std::optional<AuthoritativeItemKey> materializedItem;

Expand All @@ -58,10 +59,16 @@ struct PromptDispatch {
std::optional<std::string> expectedTurnId;
};

struct PromptVisualAlias {
LocalPromptKey key;
std::optional<AuthoritativeItemKey> admissionAnchor;
std::uint64_t admissionOrdinal = 0;
};

struct AuthoritativeItem {
AuthoritativeItemKey key;
const ItemPresentation *presentation = nullptr;
std::optional<LocalPromptKey> localPromptKey;
std::optional<PromptVisualAlias> promptAlias;
};

struct AuthoritativeItemIndex {
Expand Down Expand Up @@ -116,7 +123,7 @@ class PromptCoordinator final {
// may bind before acknowledgement so the awaiting card is never duplicated;
// the content fallback is used only after the real operation callback. Fully
// resolved submissions are removed after their accepted transition while a
// compact visual-key alias retains the admitted card identity.
// compact visual alias retains the admitted card identity and boundary.
void reconcile(const std::string &threadId,
const ThreadPresentation &authoritativeThread,
qint64 nowMilliseconds);
Expand All @@ -141,7 +148,7 @@ class PromptCoordinator final {
AuthoritativeItemIndex &authoritativeItems) const;

std::map<std::string, std::vector<PromptSubmission>> byThread;
std::map<std::string, std::map<AuthoritativeItemKey, LocalPromptKey>>
std::map<std::string, std::map<AuthoritativeItemKey, PromptVisualAlias>>
visualAliasesByThread;
std::uint64_t nextSubmissionId = 1;
std::uint64_t nextAdmissionOrdinal = 1;
Expand Down
160 changes: 160 additions & 0 deletions tests/codex/ConversationProjectionTest.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -218,6 +218,23 @@ bool testDispatchChoiceAndPreHydrationTail() {
result &=
expect(keys.size() == 3 && keys.back() == CardKey{LocalPromptKey{tailId}},
"a pre-hydration prompt stays after retained history");

PromptCoordinator recovering;
ThreadPresentation empty;
empty.id = "thread-recovering";
const auto recoveringId =
recovering.admit(empty.id, QStringLiteral("retry after resume"), {},
nlohmann::json::object(), &empty, std::nullopt, 500);
result &= expect(recovering.beginNext(empty.id).has_value() &&
recovering.requeue(empty.id, recoveringId),
"an empty-thread dispatch can return to hydration");
ThreadPresentation recovered = baseThread(empty.id);
const auto awaitingHydration = ConversationProjection::project(
recovered, recovering.submissions(empty.id), 80, 501);
result &= expect(
awaitingHydration.cardKeys().back() ==
CardKey{LocalPromptKey{recoveringId}},
"a requeued prompt returns to the unresolved retained-history tail");
return result;
}

Expand Down Expand Up @@ -266,6 +283,148 @@ bool testClientIdentityBindsBeforeAcknowledgement() {
return result;
}

bool testFirstResponseOrderIsAdmissionStable() {
ThreadPresentation reasoningFirst;
reasoningFirst.id = "thread-reasoning-first";
PromptCoordinator prompts;
const auto promptId = prompts.admit(
reasoningFirst.id, QStringLiteral("new prompt"), {},
nlohmann::json::object(), &reasoningFirst, std::nullopt, 600);
const auto dispatch = prompts.beginNext(reasoningFirst.id);
bool result =
expect(dispatch.has_value(), "an empty-thread prompt begins dispatch");
if (!dispatch)
return false;

addTurn(reasoningFirst, "turn-new");
appendItem(reasoningFirst, "turn-new",
item("reasoning", {{"type", "reasoning"},
{"summary", nlohmann::json::array()}}));
prompts.reconcile(reasoningFirst.id, reasoningFirst, 601);
const AuthoritativeItemKey reasoningKey{reasoningFirst.id, "turn-new",
"reasoning"};
const auto beforeUser = ConversationProjection::project(
reasoningFirst, prompts.submissions(reasoningFirst.id), 80, 601);
result &=
expect(beforeUser.cardKeys() ==
std::vector<CardKey>{LocalPromptKey{promptId}, reasoningKey},
"reasoning arriving first remains after its admitted prompt");

appendItem(reasoningFirst, "turn-new",
item("user-new",
{{"type", "userMessage"},
{"clientId", dispatch->clientUserMessageId},
{"content", {{{"type", "text"}, {"text", "new prompt"}}}}}));
prompts.reconcile(reasoningFirst.id, reasoningFirst, 602);
const auto materialized = ConversationProjection::project(
reasoningFirst, prompts.submissions(reasoningFirst.id), 80, 602);
result &=
expect(materialized.cardKeys() ==
std::vector<CardKey>{LocalPromptKey{promptId}, reasoningKey},
"early user-message materialization cannot invert the cards");

result &= expect(prompts.acknowledge(reasoningFirst.id, promptId,
std::string("turn-new"), 700),
"the reasoning-first prompt is acknowledged");
prompts.reconcile(reasoningFirst.id, reasoningFirst, 700);
const auto transitioning = ConversationProjection::project(
reasoningFirst, prompts.submissions(reasoningFirst.id), 80, 700);
result &=
expect(transitioning.cardKeys() ==
std::vector<CardKey>{LocalPromptKey{promptId}, reasoningKey},
"the animated-to-blue transition retains prompt order");

auto compactedItems =
indexAuthoritativeItems(reasoningFirst.id, &reasoningFirst);
prompts.reconcile(reasoningFirst.id, compactedItems, 1200);
const auto compacted = ConversationProjection::project(
compactedItems, &reasoningFirst, prompts.submissions(reasoningFirst.id),
80, 1200);
const VisibleCardData *bluePrompt = compacted.find(LocalPromptKey{promptId});
result &= expect(
compacted.cardKeys() ==
std::vector<CardKey>{LocalPromptKey{promptId}, reasoningKey} &&
bluePrompt && bluePrompt->kind == CardKind::UserMessage,
"the compact blue card retains the original admission boundary");

ThreadPresentation continued = baseThread("thread-continued");
PromptCoordinator continuedPrompts;
const auto continuedId = continuedPrompts.admit(
continued.id, QStringLiteral("continued prompt"), {},
nlohmann::json::object(), &continued, std::nullopt, 750);
const auto continuedDispatch = continuedPrompts.beginNext(continued.id);
result &= expect(continuedDispatch.has_value(),
"a continued-thread prompt begins dispatch");
if (!continuedDispatch)
return false;
addTurn(continued, "turn-continued");
appendItem(
continued, "turn-continued",
item("reasoning-continued",
{{"type", "reasoning"}, {"summary", nlohmann::json::array()}}));
appendItem(
continued, "turn-continued",
item("user-continued",
{{"type", "userMessage"},
{"clientId", continuedDispatch->clientUserMessageId},
{"content", {{{"type", "text"}, {"text", "continued prompt"}}}}}));
continuedPrompts.reconcile(continued.id, continued, 751);
result &=
expect(continuedPrompts.acknowledge(continued.id, continuedId,
std::string("turn-continued"), 800),
"the continued-thread prompt is acknowledged");
auto continuedItems = indexAuthoritativeItems(continued.id, &continued);
continuedPrompts.reconcile(continued.id, continuedItems, 1300);
const auto continuedCompacted = ConversationProjection::project(
continuedItems, &continued, continuedPrompts.submissions(continued.id),
80, 1300);
const auto continuedKeys = continuedCompacted.cardKeys();
const auto continuedPrompt =
std::ranges::find(continuedKeys, CardKey{LocalPromptKey{continuedId}});
const auto continuedReasoning = std::ranges::find(
continuedKeys,
CardKey{AuthoritativeItemKey{continued.id, "turn-continued",
"reasoning-continued"}});
result &= expect(
continuedPrompt != continuedKeys.end() &&
continuedReasoning != continuedKeys.end() &&
continuedPrompt < continuedReasoning,
"a continued-thread blue card cannot move below earlier reasoning");

ThreadPresentation userFirst;
userFirst.id = "thread-user-first";
PromptCoordinator ordinaryPrompts;
const auto ordinaryId = ordinaryPrompts.admit(
userFirst.id, QStringLiteral("ordinary prompt"), {},
nlohmann::json::object(), &userFirst, std::nullopt, 800);
const auto ordinaryDispatch = ordinaryPrompts.beginNext(userFirst.id);
result &= expect(ordinaryDispatch.has_value(),
"the user-first prompt begins dispatch");
if (!ordinaryDispatch)
return false;
addTurn(userFirst, "turn-ordinary");
appendItem(
userFirst, "turn-ordinary",
item("user-ordinary",
{{"type", "userMessage"},
{"clientId", ordinaryDispatch->clientUserMessageId},
{"content", {{{"type", "text"}, {"text", "ordinary prompt"}}}}}));
appendItem(
userFirst, "turn-ordinary",
item("reasoning-ordinary",
{{"type", "reasoning"}, {"summary", nlohmann::json::array()}}));
ordinaryPrompts.reconcile(userFirst.id, userFirst, 801);
const auto userBeforeReasoning = ConversationProjection::project(
userFirst, ordinaryPrompts.submissions(userFirst.id), 80, 801);
result &= expect(userBeforeReasoning.cardKeys() ==
std::vector<CardKey>{
LocalPromptKey{ordinaryId},
AuthoritativeItemKey{userFirst.id, "turn-ordinary",
"reasoning-ordinary"}},
"the ordinary user-first event order remains unchanged");
return result;
}

bool testAnchoredDuplicatePrompts() {
ThreadPresentation thread = baseThread("thread-duplicates");
PromptCoordinator prompts;
Expand Down Expand Up @@ -636,6 +795,7 @@ int main() {
result &= testQueueIsolationAndRealAcknowledgement();
result &= testDispatchChoiceAndPreHydrationTail();
result &= testClientIdentityBindsBeforeAcknowledgement();
result &= testFirstResponseOrderIsAdmissionStable();
result &= testAnchoredDuplicatePrompts();
result &= testCommandOutputVisibility();
result &= testUserMessageImages();
Expand Down