Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
118 commits
Select commit Hold shift + click to select a range
c2c920d
chore: files changed crates/tinyhivemind-hives/src/storage/test.rs
senamakel Oct 4, 2026
98a3cda
chore: files changed crates/tinyhivemind-hives/src/storage/test.rs,cr…
senamakel Oct 4, 2026
cb59da2
chore: files changed crates/tinyhivemind-hives/src/storage/test.rs
senamakel Oct 4, 2026
5af0e46
chore: files changed crates/tinyhivemind-hives/src/error.rs,crates/ti…
senamakel Oct 4, 2026
9f7ecd1
chore: files changed crates/tinyhivemind-hives/src/storage/mod.rs
senamakel Oct 4, 2026
94a0ddb
chore: files changed crates/tinyhivemind-hives/src/storage/sqlite.rs
senamakel Oct 4, 2026
0871534
chore: files changed crates/tinyhivemind-hives/src/coordinator/transa…
senamakel Oct 4, 2026
71516a4
chore: files changed crates/tinyhivemind-hives/src/coordinator/mod.rs…
senamakel Oct 4, 2026
feabfa0
chore: files changed crates/tinyhivemind-hives/src/coordinator/messag…
senamakel Oct 4, 2026
48b212b
chore: files changed crates/tinyhivemind-hives/src/coordinator/schedu…
senamakel Oct 4, 2026
f777fc2
chore: files changed crates/tinyhivemind-hives/src/coordinator/schedu…
senamakel Oct 4, 2026
66a98bf
chore: files changed crates/tinyhivemind-hives/src/coordinator/test/f…
senamakel Oct 4, 2026
8388705
chore: files changed crates/tinyhivemind-hives/src/coordinator/test/f…
senamakel Oct 4, 2026
b71cec3
chore: files changed crates/tinyhivemind-hives/src/coordinator/test/m…
senamakel Oct 4, 2026
65bf8ca
test(coordinator): exercise storage failure path in registration test
senamakel Oct 4, 2026
b40aa42
docs(tinyhivemind-hives): document async storage port and executor-ne…
senamakel Oct 4, 2026
1098fec
refactor(openhuman): make agent registration async
senamakel Oct 4, 2026
fd3b673
test(tools): add coverage for direct tool dispatch
senamakel Oct 4, 2026
bed78cd
test(openhuman): add host and tool integration tests
senamakel Oct 4, 2026
017d513
test(openhuman): adapt storage and api tests to async storage trait
senamakel Oct 4, 2026
50e0a68
chore(openhuman): add crate scaffold
senamakel Oct 4, 2026
bea62e9
feat(examples): add hive and openhuman example programs
senamakel Oct 4, 2026
18796c4
style: reformat openhuman examples to satisfy rustfmt
senamakel Oct 4, 2026
50c1b38
fix(proof): await call count load in dynamic proof
senamakel Oct 4, 2026
258fce6
feat(coordinator): expose episode status to hosts
senamakel Oct 4, 2026
3b6854b
chore(coordinator): add debug logging to scheduler loop
senamakel Oct 4, 2026
6337476
chore(scheduler): add temporary debug logging
senamakel Oct 4, 2026
e338188
chore(coordinator): remove debug prints and cover stalled episodes
senamakel Oct 4, 2026
b63822d
fix(coordinator): checkpoint stalled episodes before marking them fin…
senamakel Oct 4, 2026
fd7f072
test(coordinator): reuse setup helper in episode snapshot test
senamakel Oct 4, 2026
b60154d
test(coordinator): drive episode snapshot cases from per-agent handlers
senamakel Oct 4, 2026
d89af49
feat(coordinator): add starters to choose who begins a hive episode
senamakel Oct 4, 2026
dce6e0a
test(openhuman): fix struct literal syntax in continuity test helper
senamakel Oct 4, 2026
097ca17
feat(coordinator): attach host note to released agent's next turn
senamakel Oct 4, 2026
14fccac
test(openhuman): cover resumption note rendering in turn prompts
senamakel Oct 4, 2026
367d052
feat(openhuman): prepend host resumption note to turn prompt
senamakel Oct 4, 2026
28dd74d
test(openhuman): add host hook and type tests
senamakel Oct 4, 2026
602d2f4
feat(host): pass turn scope and options to hooks
senamakel Oct 4, 2026
c4377e4
feat(host): add replace_agent to rebuild an agent handle
senamakel Oct 4, 2026
48aa4c3
feat(openhuman): gate outbound send tools behind a host policy
senamakel Oct 4, 2026
50686b8
docs(adr): index ADR 0031 in the ADR README
senamakel Oct 4, 2026
3d9e2c3
docs(opencompany-migration): document async storage and host seams
senamakel Oct 4, 2026
eced472
docs: document async storage, host seams, and new coordinator APIs
senamakel Oct 4, 2026
c1d1af6
Merge branch 'land-hive-lab' into hives-host-seams
senamakel Oct 4, 2026
20e410f
Merge remote-tracking branch 'upstream/main' into hives-host-seams
senamakel Oct 4, 2026
90f3910
chore: files changed docs/adr/0031-expose-host-seams-over-async-incre…
senamakel Oct 4, 2026
f21a058
docs: remove adr 31 on coordinator host seams
senamakel Oct 4, 2026
5108124
docs(adr): add index for architecture decision records
senamakel Oct 4, 2026
eaa98c8
docs: add opencompany migration guide
senamakel Oct 4, 2026
43dfd7e
docs: add dynamic hives spec
senamakel Oct 4, 2026
f5df21b
refactor(coordinator): extract peer registration into helper
senamakel Oct 4, 2026
e6f9cd1
chore(coordinator): remove unused transaction helper
senamakel Oct 4, 2026
9621f9f
chore: files changed crates/tinyhivemind-hives/src/coordinator/transa…
senamakel Oct 4, 2026
035abc0
chore(storage): add missing trailing newline to types.rs
senamakel Oct 4, 2026
841bf28
refactor(storage): derive Debug for storage types
senamakel Oct 4, 2026
b1bcbe3
test(coordinator): add transaction tests for hive coordinator
senamakel Oct 4, 2026
83852d5
chore(test): add transaction coordinator tests
senamakel Oct 4, 2026
78e1568
chore(test): add transaction coordinator tests
senamakel Oct 4, 2026
11101ae
test(coordinator): add transaction tests for hive coordinator
senamakel Oct 4, 2026
c916dc9
chore(test): add storage tests
senamakel Oct 4, 2026
338e29a
chore(test): add storage tests
senamakel Oct 4, 2026
6282d30
chore(test): add transaction coordinator tests
senamakel Oct 4, 2026
86ce6d2
test(coordinator): add transaction tests for hive coordination
senamakel Oct 4, 2026
4a76b6a
test(coordinator): add transaction tests for hive coordinator
senamakel Oct 4, 2026
043ac59
refactor(storage): rename HivePartition to HivePartitionKey
senamakel Oct 4, 2026
0fcb232
fix(coordinator): retry transaction commit on transient failures
senamakel Oct 4, 2026
cce2124
style(coordinator): reflow should_reapply check
senamakel Oct 4, 2026
882522d
refactor(storage): derive default for hive storage types
senamakel Oct 4, 2026
a3580f7
fix(error): add missing error variants for hive operations
senamakel Oct 4, 2026
b919862
refactor(coordinator): extract peer registration into helper
senamakel Oct 4, 2026
08916aa
refactor(coordinator): extract peer registration into helper
senamakel Oct 4, 2026
08e5678
refactor(coordinator): simplify transaction handling
senamakel Oct 4, 2026
1611c3c
refactor(coordinator): simplify transaction handling
senamakel Oct 4, 2026
e395965
refactor(coordinator): simplify transaction handling
senamakel Oct 4, 2026
6f1cd50
refactor(coordinator): simplify transaction handling
senamakel Oct 4, 2026
f32e67d
refactor(coordinator): simplify transaction handling
senamakel Oct 4, 2026
297caa1
refactor(coordinator): simplify transaction handling
senamakel Oct 4, 2026
594a56c
refactor(coordinator): simplify transaction handling
senamakel Oct 4, 2026
a3d4cb3
refactor(storage): derive default for hive storage types
senamakel Oct 4, 2026
850fc5e
test(coordinator): add transaction tests for hive coordinator
senamakel Oct 4, 2026
9cc144e
test(coordinator): add transaction tests for hive coordinator
senamakel Oct 4, 2026
c248876
chore(test): remove unused transaction test helpers
senamakel Oct 4, 2026
52ebed8
test(coordinator): add transaction tests for hive coordinator
senamakel Oct 4, 2026
2c1ac7f
refactor(coordinator): simplify transaction handling
senamakel Oct 4, 2026
b4bfab5
refactor(coordinator): restructure transaction handling
senamakel Oct 4, 2026
5af6bf8
chore(coordinator): remove unused import
senamakel Oct 4, 2026
fec1211
test(coordinator): add transaction tests for hive coordinator
senamakel Oct 4, 2026
616959e
test(coordinator): add transaction tests for hive coordinator
senamakel Oct 4, 2026
fe22d30
chore(test): add transaction coordinator tests
senamakel Oct 4, 2026
8ca2119
chore(test): add transaction coordinator tests
senamakel Oct 4, 2026
712ac5a
test(coordinator): add transaction tests for hive coordinator
senamakel Oct 4, 2026
791edc0
test(coordinator): remove redundant transaction tests
senamakel Oct 4, 2026
5506099
test(coordinator): remove duplicate tokio test attribute
senamakel Oct 4, 2026
965c866
chore(test): remove unused transaction test helpers
senamakel Oct 4, 2026
65ccbc8
refactor(coordinator): extract transaction commit into helper
senamakel Oct 4, 2026
1ffaea9
refactor(coordinator): simplify transaction handling
senamakel Oct 4, 2026
4828a73
refactor(coordinator): extract peer registration into helper
senamakel Oct 4, 2026
a4b44d6
refactor(scheduler): extract peer selection into a helper
senamakel Oct 4, 2026
8892987
fix(coordinator): handle missing message recipients gracefully
senamakel Oct 4, 2026
15fb66d
docs(adr): document writer fencing and retention changes
senamakel Oct 4, 2026
f48dd8d
docs: document single-writer fencing and bounded inbox guarantees
senamakel Oct 4, 2026
59a0150
test: add transaction and storage coverage
senamakel Oct 4, 2026
f87aeb3
test: add transaction and storage coverage
senamakel Oct 4, 2026
3d20dde
test(coordinator): simplify interrupted record setup
senamakel Oct 4, 2026
ac4cc26
test(coordinator): relax interruption count assertion
senamakel Oct 4, 2026
c68ce50
docs(coordinator): update transactions test module summary
senamakel Oct 4, 2026
aad1b47
refactor(coordinator): remove unused conflict retry plumbing
senamakel Oct 4, 2026
7cba88e
refactor(coordinator): split transaction handling into its own module
senamakel Oct 4, 2026
fc9a6bb
chore(hives): remove unused coordinator module
senamakel Oct 4, 2026
260bc52
docs(coordinator): clarify writer epoch and deferred interruption com…
senamakel Oct 4, 2026
816c867
chore(coordinator): import AtomicU64
senamakel Oct 4, 2026
19cd10f
test(coordinator): add transaction tests for hive coordinator
senamakel Oct 4, 2026
cf76175
style(coordinator): reformat writer epoch increment and add missing f…
senamakel Oct 4, 2026
9d4c534
test(coordinator): match any ok value when asserting fencing
senamakel Oct 4, 2026
f6846e6
test(coordinator): reorder hive setup before stale coordinator
senamakel Oct 4, 2026
9b39530
test(openhuman): reject registration after coordinator construction
senamakel Oct 4, 2026
973df40
refactor(storage): derive default for storage types
senamakel Oct 4, 2026
80d8d23
test(storage): assert legacy snapshots migrate without a writer epoch
senamakel Oct 4, 2026
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
12 changes: 11 additions & 1 deletion crates/tinyhivemind-hives/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,11 @@ a desk have one identity. Each registered agent has one continuing session and
one globally serialized turn stream, even when it belongs to several hives.

Use `Coordinator` with `MemoryStorage`, default-feature `SqliteStorage`, or a
host implementation of `Storage`. The host supplies `AgentRunner` handles from
host implementation of the async `Storage` port. A commit writes one bounded
state row and appends only the new transcript rows; `RetentionPolicy` bounds
settled episodes and acknowledged deliveries. Mutating APIs are `async` and
executor-neutral; commits run outside the live lock and reload on a conflict
with another process. The host supplies `AgentRunner` handles from
one runtime; this crate never constructs agents or serializes live handles.
The OpenHuman-specific boundary lives in `tinyhivemind-openhuman`.

Expand All @@ -15,6 +19,12 @@ claims it. Hosts can create empty hives and join/leave registered agents while
the scheduler runs; unstarted removed seats are retired before claiming.
An active turn retains its captured membership.

The host reads everything with `read_transcript`, including replies to
`send_as_host`, and watches `subscribe()` (the committed revision) and
`episodes()` for settlement. `SendMessage::starters` chooses who starts an
episode without hiding the message, and `release_with` hands a parked agent
a note on its next turn.

Direct sends return durable receipts without awaiting peers. `read_direct`
exposes only the caller/peer pair, including returned replies, and starts no
turn. Hive reads enforce membership and private thread audiences.
Expand Down
22 changes: 17 additions & 5 deletions crates/tinyhivemind-hives/src/coordinator/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -2,17 +2,29 @@

| File | Responsibility |
| --- | --- |
| `mod.rs` | Shared handle, dynamic APIs, authorization snapshots, and transactions |
| `mod.rs` | Shared handle, dynamic APIs, and authorization snapshots |
| `transaction.rs` | Writer gate, incremental commits outside the live lock, conflict reload |
| `types.rs` | Host runner port and public payloads |
| `messaging.rs` | Atomic acceptance, retry IDs, attribution, and private reads |
| `messaging.rs` | Atomic acceptance, retry IDs, starters, attribution, and private reads |
| `observe.rs` | Host transcript, committed-revision watch, episode status |
| `conduct.rs` | Actual `CompletionDriver`/`Conductor` checkpoint lifecycle |
| `scheduler.rs` | FIFO reservations, concurrent futures, shutdown, cancellation |
| `test/` | Deterministic contract fixtures and behavior tests |

All clones of `Coordinator` share one scheduler and state. Transactions clone
state and replace it only after the storage CAS succeeds. Runner futures execute
outside the state lock. Conductor folding uses copied snapshots and retries if a
concurrent synchronous API changed the revision; no external tool is repeated.
state under the live lock, release it, and persist the copy under an async
writer gate. The storage commit carries the state row plus only the transcript
rows appended since the base, and the copy is published only after the CAS
succeeds. Reads never wait on storage. A revision conflict, which means another
process wrote the store, reloads and recomputes up to four times. Runner
futures execute outside every lock. A dropped drain applies its interruptions
to live state immediately, and the next commit persists them. No external tool
is repeated.

`release_with(agent, note)` stores a note on the agent record. The next claim
moves it into `TurnRequest::resumption`, once. Host-chosen `starters` must be
distinct readers of the hive message; they open the episode while every reader
still sees the message.

Work is ordered by accepted sequence, then agent ID within a round. Shared
agents retain their queued positions and never run two turns concurrently.
Expand Down
8 changes: 7 additions & 1 deletion crates/tinyhivemind-hives/src/coordinator/conduct.rs
Original file line number Diff line number Diff line change
Expand Up @@ -253,13 +253,19 @@ pub(super) async fn prepare(state: &mut StoredState, options: &CoordinatorOption
}
// An empty parked wave must wait for an explicit release.
episode.waiting = episode.pending.is_empty() && !conductor.parked().is_empty();
checkpoint(&conductor, &mut episode)?;
}
Err(error) => {
// The conductor still reports itself unfinished after a stall,
// so settle after checkpointing; otherwise the episode is
// re-prepared, and fails again, on every pass.
checkpoint(&conductor, &mut episode)?;
episode.finished = true;
episode.failure = Some(error.to_string());
episode.pending.clear();
episode.wave_open = false;
}
}
checkpoint(&conductor, &mut episode)?;
state.episodes[index] = episode;
}
Ok(())
Expand Down
50 changes: 40 additions & 10 deletions crates/tinyhivemind-hives/src/coordinator/messaging.rs
Original file line number Diff line number Diff line change
Expand Up @@ -9,25 +9,26 @@ impl Coordinator {
/// # Errors
/// Returns unknown sender/destination, missing membership, invalid thread,
/// conflicting retry identity or persistence errors.
pub fn send(&self, request: SendMessage) -> Result<Receipt> {
pub async fn send(&self, request: SendMessage) -> Result<Receipt> {
if request.sender == HOST_ID {
return Err(Error::InvalidIdentifier("sender"));
}
self.accept(request)
self.accept(request).await
}
/// Submit through the reserved host identity rather than impersonating an agent.
/// The request's sender field is overwritten.
/// # Errors
/// Returns destination, visibility, retry or storage validation errors.
pub fn send_as_host(&self, mut request: SendMessage) -> Result<Receipt> {
pub async fn send_as_host(&self, mut request: SendMessage) -> Result<Receipt> {
request.sender = HOST_ID.into();
self.accept(request)
self.accept(request).await
}
fn accept(&self, request: SendMessage) -> Result<Receipt> {
async fn accept(&self, request: SendMessage) -> Result<Receipt> {
identifier(&request.message_id, "message id")?;
if request.message_id.starts_with("hivemind:") {
return Err(Error::InvalidIdentifier("reserved message id"));
}
let retention = self.inner.options.retention;
self.update(|state| {
if let Some(old) = state.accepted.get(&request.message_id) {
if old != &request {
Expand Down Expand Up @@ -67,6 +68,9 @@ impl Coordinator {
};
match &request.destination {
Destination::Agent(_) => {
for agent_id in &recipients {
retention.admit_pending(state, agent_id)?;
}
for agent_id in recipients {
state.deliveries.push(Delivery {
sequence,
Expand Down Expand Up @@ -94,7 +98,11 @@ impl Coordinator {
hive,
opened_at: sequence,
thread: request.thread,
starters: recipients,
starters: if request.starters.is_empty() {
recipients.clone()
} else {
request.starters.clone()
},
conductor: None,
pending: Vec::new(),
wave_open: false,
Expand All @@ -109,10 +117,11 @@ impl Coordinator {
.accepted
.insert(request.message_id.clone(), request.clone());
Ok(Receipt {
message_id: request.message_id,
message_id: request.message_id.clone(),
sequence,
})
})
.await
}
/// Read the caller's direct conversation with a registered peer in durable
/// sequence order, including replies returned by the peer's runner.
Expand Down Expand Up @@ -220,7 +229,10 @@ fn recipients(state: &StoredState, request: &SendMessage) -> Result<Vec<String>>
Ok(match &request.destination {
Destination::Agent(id) => {
known_agent(state, id)?;
if request.thread.is_some() || !request.only_for.is_empty() {
if request.thread.is_some()
|| !request.only_for.is_empty()
|| !request.starters.is_empty()
{
return Err(Error::InvalidIdentifier("direct message attribution"));
}
vec![id.clone()]
Expand Down Expand Up @@ -268,13 +280,31 @@ fn recipients(state: &StoredState, request: &SendMessage) -> Result<Vec<String>>
{
return Err(Error::InvalidThread(request.thread.unwrap_or_default()));
}
selected
let selected: Vec<_> = selected
.into_iter()
.filter(|recipient| scope.is_empty() || scope.contains(recipient))
.collect()
.collect();
starters_among(&request.starters, &selected, id)?;
selected
}
})
}
/// Host-chosen starters must be distinct readers of the message.
fn starters_among(starters: &[String], readers: &[String], hive_id: &str) -> Result<()> {
let mut unique = BTreeSet::new();
for starter in starters {
if !readers.contains(starter) {
return Err(Error::NotMember {
agent_id: starter.clone(),
hive_id: hive_id.into(),
});
}
if !unique.insert(starter) {
return Err(Error::DuplicateMember(starter.clone()));
}
}
Ok(())
}

/// Root-private threads can only be read by their original participants.
pub(super) fn visible_in(state: &StoredState, message: &Message, agent_id: &str) -> bool {
Expand Down
Loading
Loading