We can assume all the processors that are relevant are online; no processor core will try to invalidate an item while being cut off from communication, which is then re-established later. (If a processor comes online dynamically, we can assume it's coming in with a clear cache.)
Bus snooping is possible: every cache can see every memory access, and based on that, it can invalidate (or even update) entries that would be made stale due to stores. In a distributed system, no such thing is possible; you can't feasibly have 10,000 geographically distributed nodes all watching every transaction.
We can also assume that any remaining race conditions are handled by the software: that the software will use tools like memory barriers in the right places; we don't have to solve every possible concurrency problem that interacts with caching effects.