Skip to main content

Federated, secure, and AI-optimized distributed computing platform

Project description

Flock

Brokerless, decentralised, peer-to-peer distributed task execution framework for Python.

Flock eliminates all centralised brokers (Redis, RabbitMQ, Kafka, Celery, etc.) by implementing a complete distributed computing stack directly in Python: from TCP transport and peer discovery through to Raft consensus, distributed scheduling, worker execution, result collection, and automatic failure recovery.


Features

  • Fully brokerless – no external infrastructure dependencies.
  • Peer-to-peer networking – asyncio TCP transport with binary packet framing.
  • Dynamic peer discovery – announce/query/expire lifecycle over the MessageBus.
  • Cluster membership – join/ack handshake with versioned snapshot synchronisation.
  • Heartbeat & failure detection – periodic health monitoring with configurable timeouts.
  • Distributed task scheduler – priority queues, constraint validation, EventBus notifications.
  • Intelligent task placement – capability matching, resource-aware ranking, assignment handshakes.
  • Worker runtime – pluggable Thread, Process, and Async execution backends.
  • Distributed result collection – checksum-verified payloads, async Future completion.
  • Retry & recovery engine – Fixed/Linear/Exponential-Jitter backoff, automatic reassignment.
  • Raft consensus – leader election with randomised timeouts, AppendEntries log replication, quorum-gated commit index, and strict log completeness enforcement.
  • Transport-independent – all subsystems communicate through the MessageBus; no subsystem depends directly on networking code.
  • mypy strict – entire codebase verified at --strict level.

Architecture

┌─────────────────────────────────────────────────────┐
│                   Application Layer                  │
├──────────────────┬──────────────────────────────────┤
│  ConsensusService│  Milestone E: Reliability         │
│  (Raft Leader    │  RecoveryEngine (Phase 11)        │
│  Election +      │  RetryPolicyEngine (Phase 11)     │
│  Log Replication)│                                   │
├──────────────────┴──────────────────────────────────┤
│              Milestone D: Distributed Execution      │
│  WorkerRuntimeService  │  ResultService              │
│  PlacementEngine        │  ResultRegistry            │
├─────────────────────────────────────────────────────┤
│              Milestone C: Distributed Scheduling     │
│  TaskSchedulerService   │  PlacementEngine           │
├─────────────────────────────────────────────────────┤
│              Milestone B: Cluster Formation          │
│  HeartbeatService       │  ClusterMembershipService  │
│  HealthRegistry          │  MembershipRegistry       │
├─────────────────────────────────────────────────────┤
│              Milestone A: Core Infrastructure        │
│  DiscoveryService       │  MessageBus / EventBus     │
│  TcpTransport           │  JsonSerializer / Packet   │
└─────────────────────────────────────────────────────┘

Project Status

Phase Name Status
1 Setup and Architecture Core ✓ Complete
2 Protocol Serialization & Transport Layer ✓ Complete
3 Communication & Messaging Core ✓ Complete
4 Peer Discovery ✓ Complete
5 Cluster Membership ✓ Complete
6 Heartbeat & Failure Detection ✓ Complete
7 Distributed Task Scheduler ✓ Complete
8 Distributed Task Placement Engine ✓ Complete
9 Worker Runtime & Execution Engine ✓ Complete
10 Distributed Result Collection & Completion Pipeline ✓ Complete
11 Distributed Retry & Recovery Engine ✓ Complete
12 Distributed Raft Consensus Engine ✓ Complete
13 Persistent Distributed Log & Snapshot Management ⬜ Planned

Current Version: 0.6.0
Total Tests: 138 passing (0 failing, 0 skipped)
mypy: strict / 0 issues


Package Structure

src/flock/
├── core/           # Common primitives, types, exceptions
├── protocol/       # Binary packet framing, MessageType constants (53 defined)
├── serialization/  # JSON / MessagePack adapters
├── transport/      # asyncio TCP transport
├── messaging/      # MessageBus, MessageRouter, middleware, EventBus
├── discovery/      # Peer discovery service and registry
├── cluster/        # Cluster membership service and registry
├── heartbeat/      # Heartbeat service and health registry
├── scheduler/      # Distributed task scheduler
├── placement/      # Task placement engine
├── runtime/        # Worker runtime and execution backends
├── results/        # Result collection and completion pipeline
├── recovery/       # Retry policy engine and recovery coordinator
└── consensus/      # Raft consensus (Phase 12)
    ├── exceptions.py   – 7 typed exception classes
    ├── models.py       – 13 immutable Pydantic models
    ├── log.py          – ConsensusLog (thread-safe, 1-based)
    ├── state_machine.py– RaftStateMachine (role FSM)
    ├── election.py     – ElectionEngine (vote solicitation)
    ├── replication.py  – ReplicationEngine (AppendEntries)
    └── service.py      – ConsensusService (orchestrator)

Raft Consensus (Phase 12)

Leader Election

Nodes start as followers with a randomised election timeout (150–300ms). On timeout, a follower:

  1. Increments its term and becomes a candidate.
  2. Votes for itself.
  3. Broadcasts RAFT_REQUEST_VOTE to all peers.
  4. Becomes leader on receiving a majority of votes.

Vote granting follows Raft §5.2 + §5.4.1: a voter only grants its vote when the candidate's term is current and its log is at least as up-to-date as the voter's log (prevents losing committed entries on leadership change).

Log Replication

The leader periodically sends RAFT_APPEND_ENTRIES to all followers. Followers apply the 5-step Raft receiver algorithm, truncating conflicting entries and appending new ones. The commit index advances only after a quorum acknowledges replication.

EventBus Events

Event Payload
consensus.leader.elected {leader_id, term}
consensus.term.changed {old_term, new_term}
consensus.log.committed {index, entry_id, term}
consensus.replication.failed {peer_id, error}

Usage

from flock.consensus.service import ConsensusService

service = ConsensusService(
    node_id="node-1",
    message_bus=bus,
    event_bus=events,
    membership_registry=registry,
    min_election_timeout=0.15,
    max_election_timeout=0.30,
    heartbeat_interval=0.05,
)

await service.start()

# Submit a command (leader only)
if service.is_leader():
    entry = await service.submit_command(b"my-payload")

Installation

pip install -e .

Requirements: Python ≥ 3.11, pydantic ≥ 2.0, structlog, msgpack.


Running Tests

python -m pytest tests/ -v

Documentation

  • docs/adr/ — Architecture Decision Records (0001–0012)
  • docs/audits/ — Phase audit reports and retrospectives
  • CHANGELOG.md — Full version history
  • PROJECT_STATE.json — Machine-readable project state

Contributing

This is a research/educational project demonstrating how to build a complete distributed task execution framework from first principles. Contributions, suggestions, and questions are welcome.


License

MIT License. See LICENSE for details.

Project details


Download files

Download the file for your platform. If you're not sure which to choose, learn more about installing packages.

Source Distribution

flock_p2p-1.0.0.tar.gz (336.6 kB view details)

Uploaded Source

Built Distribution

If you're not sure about the file name format, learn more about wheel file names.

flock_p2p-1.0.0-py3-none-any.whl (385.7 kB view details)

Uploaded Python 3

File details

Details for the file flock_p2p-1.0.0.tar.gz.

File metadata

  • Download URL: flock_p2p-1.0.0.tar.gz
  • Upload date:
  • Size: 336.6 kB
  • Tags: Source
  • Uploaded using Trusted Publishing? No
  • Uploaded via: twine/6.2.0 CPython/3.11.4

File hashes

Hashes for flock_p2p-1.0.0.tar.gz
Algorithm Hash digest
SHA256 1d7ce080dc7c6898795114e5b662dfb094c354f016f39944039396e00fe5d589
MD5 ce93840035b76d441506f825ce5284c1
BLAKE2b-256 a501616e25289bca54c32ed649ef4c1200d234674aaeadd7612a6a9dc77e966a

See more details on using hashes here.

File details

Details for the file flock_p2p-1.0.0-py3-none-any.whl.

File metadata

  • Download URL: flock_p2p-1.0.0-py3-none-any.whl
  • Upload date:
  • Size: 385.7 kB
  • Tags: Python 3
  • Uploaded using Trusted Publishing? No
  • Uploaded via: twine/6.2.0 CPython/3.11.4

File hashes

Hashes for flock_p2p-1.0.0-py3-none-any.whl
Algorithm Hash digest
SHA256 bc8f4e479758de396ec35815d4e7d3e26b447b5ff383d4cbaecf94bf06b4e5d0
MD5 24c43f85a0d42336e67da456b9afcf4b
BLAKE2b-256 2f19e32ca22fcad761ad1a6492c9c45617cd453e8adc3b762c9054373ac569dd

See more details on using hashes here.

Supported by

AWS Cloud computing and Security Sponsor Datadog Monitoring Depot Continuous Integration Fastly CDN Google Download Analytics Pingdom Monitoring Sentry Error logging StatusPage Status page